Aiuto per il colloquio
Preparazione al colloquio
Riprendi
Ricerca di lavoro
Risorse
Abbonamento
IT
Telegram
Accedi
Apri il menu
Aiuto per il colloquio
Assistente AI
Estensione del browser
Risorse
Scarica
Domande frequenti/Aiuto
Preparazione al colloquio
Banca delle domande
Banca delle interviste
Interviste finte
Mentori
Risorse
Domande frequenti/Aiuto
Ricerca di lavoro
Automazione della ricerca di lavoro
Analisi di mercato
Offerte di lavoro da Telegram
Risorse
Lavora con noi
Blog
Domande frequenti/documentazione
Telegram Blog
Telegram Supporto
Legale
Rimborso
Termini
Privacy
Titolari dei diritti
Torniamo alle domande
Banca delle domande
Data Engineer - Sviluppatore
Domanda
Middle
da interview.utils import get_clickhouse_client from airflow import DAG from airflow.operators.python import PythonOperator from airflow.sensors.external_task import ExternalTaskSensor from datetime import datetime import pandas as pd import clickhouse_driver import os CLICKHOUSE_CLIENT = get_clickhouse_client() default_args = { "start_date": datetime(2024, 1, 1) } with DAG( dag_id="datamarts.daily_revenue_per_country", default_args=default_args, schedule_interval="@daily", catchup=False ) as dag: transactions_sensor = S3KeySensor( task_id="transactions_sensor", bucket_key="data/transactions_{}.csv".format(datetime.now().strftime("%Y-%m-%d")), bucket_name="my-bucket", aws_conn_id="aws_default", timeout=600, poke_interval=30, mode="poke" ) def extract_from_s3(**kwargs): df = pd.read_csv("s3://my-bucket/data/transactions_{}.csv".format(datetime.now().strftime("%Y-%m-%d"))) kwargs["ti"].xcom_push(key="df", value=df.to_dict()) def load_to_raw_table(**kwargs): df = pd.DataFrame(kwargs["ti"].xcom_pull(task_ids="extract", key="df")) rows = [tuple(r) for r in df[["transaction_id", "user_id", "amount", "created_at"]].to_numpy()] CLICKHOUSE_CLIENT.execute( ... )
Contrassegna come completato
Precedente
Avanti