Data Engineer
Cuéntame sobre ti: experiencia, tareas, características o logros de los que estés orgulloso.
¿Tienes experiencia con FastAPI?
¿Cuánto tiempo en total trabajó con Spark y qué tan profundo se involucró en él?
dict1 = { (1, 2), [3, 4, 5] : 0 } var = 1, 2
Cuéntame sobre un problema interesante en la empresa en los últimos 2,5 años, por ejemplo, relacionado con ClickHouse.
¿En qué situación la búsqueda en un diccionario puede degradarse hasta el peor caso?
¿Cómo abordó los problemas de optimización cuando la consulta tarda mucho y realiza un escaneo completo? ¿Cómo abordaría una tarea similar nueva?
¿Cómo agregar una media móvil de la suma de pedidos de los días actual y dos anteriores por usuario?
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
¿Qué hace el parámetro depends_on_past en Airflow?
¿Para qué necesitabas Data Vault en general?
Cuéntame sobre la tarea de configuración de replicación que resolviste.
¿Se utilizó el método de comparación de hashes de registros para calcular la diferencia entre la fuente y la tabla de destino?
¿Qué ocurría cuando se activaba la restricción de calidad de datos: una alerta, una caída, un reinicio?
Compare los enfoques ETL y ELT: ¿cuál es la diferencia y cuándo se aplica cada uno?
¿Cómo trabajaban con 1C: extraían datos directamente de la base de datos, a través de un bus o de otra manera?
¿Con qué bases de datos ha trabajado? ¿Cómo funciona MongoDB y cómo se escala?
¿Cómo determinaron que la entrada ya existe y no es necesario volver a grabarla?
NULL más 5, ¿cuánto será?
¿Se puede leer datos en paralelo desde un solo archivo Parquet en Spark, o solo con un núcleo en una tarea?