iz interview.utils uvozite get_clickhouse_client
iz airflow uvozite DAG
iz airflow.operators.python uvozite PythonOperator
iz airflow.sensors.external_task uvozite ExternalTaskSensor
iz datetime uvozite datetime
uvoz pandas kot pd
uvoz clickhouse_driver
uvoz os
CLICKHOUSE_CLIENT = get_clickhouse_client()
default_args = {
"start_date": datetime(2024, 1, 1)
}
z DAG(
dag_id="datamarts.daily_revenue_per_country",
default_args=default_args,
schedule_interval="@daily",
catchup=False
) kot 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( ... )