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
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( "INSERT INTO raw.transactions (transaction_id, user_id, amount, created_at) VALUES", rows ) def build_aggregate_view(): query = """ INSERT INTO datamarts.daily_revenue_per_country SELECT toDate(r.created_at) as event_date, u.country, sum(r.amount) as total_revenue FROM raw.transactions r LEFT JOIN core.userMetadata u ON r.user_id = u.UserId WHERE toDate(r.created_at) = '{}' GROUP BY event_date, u.country """.format(datetime.now().strftime("%Y-%m-%d")) CLICKHOUSE_CLIENT.execute(query) 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" ) extract = PythonOperator( task_id="extract", python_callable=extract_from_s3, provide_context=True ) load = PythonOperator( task_id="load", python_callable=load_to_raw_table, provide_context=True ) aggregate = PythonOperator( task_id="aggregate", python_callable=build_aggregate_view ) transactions_sensor >> extract >> load >> aggregate
Contrassegna come completato
Precedente
Avanti