Ayuda para la entrevista
Preparación de la entrevista
Currículum
Búsqueda de empleo
Recursos
Suscripción
ES
Telegram
Iniciar sesión
abrir menú
Ayuda para la entrevista
Asistente de IA
Extensión del navegador
Recursos
Descargar
Preguntas frecuentes / Ayuda
Preparación de la entrevista
banco de preguntas
Banco de entrevistas
Entrevistas simuladas
Mentores
Recursos
Preguntas frecuentes / Ayuda
Búsqueda de empleo
Automatización de búsqueda de empleo
Análisis de mercado
Empleos de Telegram
Recursos
Trabaja con nosotros
Blog
Preguntas frecuentes / Documentación
Telegram Blog
Telegram Soporte
Legales
Reembolso
Términos
Privacidad
Titulares de derechos
Volver a preguntas
banco de preguntas
Data Engineer - Desarrollador
Pregunta
Middle
desde 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( ... )
Marcar como completado
Anterior
Siguiente