Data Engineer
Come aggiungere una media mobile della somma degli ordini del giorno corrente e dei due giorni precedenti per utente?
Come funziona l'hash join per tabelle con 1.000 e 1.000.000 di righe?
A cosa serviva in generale Data Vault?
È stato utilizzato il metodo di confronto degli hash delle registrazioni per calcolare la differenza tra la sorgente e la tabella di destinazione?
Hai mai riscontrato l'errore ClickHouse `Too many parts` durante un inserimento? Come l'hai risolto?
Con quale versione di Airflow hai lavorato?
Con quali formati di file hai lavorato?
Quali tipi di JOIN fisici esistono in Spark?
Cosa succede fisicamente durante un INSERT, UPDATE e DELETE in Greenplum; perché dopo un DELETE lo spazio su disco non viene liberato e in cosa si differenzia DELETE da TRUNCATE?
Perché i dati possono andare persi? Fornisci alcune ragioni.
Perché stai considerando nuove offerte ora?
Parlami di te: esperienza, compiti, funzionalità o risultati di cui sei orgoglioso.
Hai esperienza con FastAPI?
Come hai lavorato con i file al riavvio del DAG dopo il suo fallimento, per non sovrascrivere i dati incompleti in S3?
In quale situazione la ricerca in un dizionario può degradarsi fino al caso peggiore?
Come hai risolto i problemi di ottimizzazione quando la query richiede molto tempo e esegue una scansione completa? Come affronteresti un nuovo compito simile?
Cos'è la partizione dei dati e come aiuta a evitare una scansione completa della tabella?
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
Cosa fa il parametro depends_on_past in Airflow?
Qual è la differenza tra ACID e il teorema CAP?