Pomoc w rozmowie kwalifikacyjnej
Przygotowanie do rozmowy kwalifikacyjnej
Wznów
Wyszukiwanie pracy
Zasoby
Subskrypcja
PL
Telegram
Zaloguj się
Otwórz menu
Pomoc w rozmowie kwalifikacyjnej
AI Assistant
Rozszerzenie przeglądarki
Zasoby
Pobierz
FAQ / Pomoc
Przygotowanie do rozmowy kwalifikacyjnej
Bank pytań
Wywiad Bank
Próbne wywiady
Mentorzy
Zasoby
FAQ / Pomoc
Wyszukiwanie pracy
Automatyzacja wyszukiwania pracy
Analiza rynku
Oferty pracy z Telegrama
Zasoby
Kariera
Blog
FAQ / Dokumentacja
Telegram Blog
Telegram Wsparcie
Prawo
Zwrot
Warunki
Prywatność
Prawo posiadacze
Wróć do pytania
Bank pytań
Data Engineer - Programista
Pytanie
Middle
z 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( ... )
Oznacz jako ukończone
Poprzednia
Dalej