Data Engineer
def extract_from_s3(**kwargs): df = pd.read_csv("s3://manobucket/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
Kāda ir atšķirība starp ACID un CAP teoremu?
Kas notiks, ja piemērosiet `s ** 2` sarakstam `s = [1, 2, 3]`?
Salīdziniet ETL un ELT pieejas: kāda ir atšķirība un kad katra no tām tiek piemērota?
Kā strādāja ar 1C: tieši no datu bāzes ņēma datus, vai caur shēnu vai citādi?
Ar ar ko datu bāzēm esat strādājis? Kā darbojas MongoDB un kā tā mērogojas?
Kā noteicāt, ka ieraksts jau pastāv un to atkārtoti ierakstīt nav nepieciešams?
NULL plus 5 — cik tas būs?
Vai Spark var lasīt datus paralēli no viena Parquet faila, vai tikai ar vienu kodolu vienā uzdevumā?
Kādās laukos notika datu apvienošana un kā tika izņemti dublikāti no vitrina?
Ko jūs zināt par spill Spark?
Kas ir RDD Apache Spark un kā tas atšķiras no citām Spark abstrakcijām?
Ar kurām Python bibliotēkām jūs strādājāt?
Kas ir CROSS JOIN un kāds ir tā pielietošanas rezultāts?
Kā jūs pievienojāties OpenMetadata un ko ar to darījāt?
Vai ir iespējams izveidot un piemērot dekoratoru bez @ sintakses?
Kādas datu sadales stratēģijas pastāv?
Ja DAG grafiks ir @daily, kad tas faktiski sākas?
Kādi ir partīciju veidi?
Kā tika tehniski īstenots liela tabulas paralēlais izkraušana no PostgreSQL uz Spark gadījumā [konteksts]?