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