Data Engineer
Pentru ce ai avut nevoie în general de Data Vault?
Pe ecranul afișat există două structuri de date; în care dintre ele căutarea valorii 2 va fi mai rapidă și de ce?
Ai experiență cu Table Engine Join în ClickHouse?
Cum se evaluează configurația pentru transferul unui terabyte de date de la PostgreSQL la ClickHouse?
Ce este RDD în Apache Spark și în ce se deosebește de alte abstracții Spark?
Ce grupuri de operatori SQL cunoști? La ce aparțin?
din 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( ... )
Rezultatul final: | customer_id | gap_days | |-------------|----------| | 1 | 123 | | 2 | 120 | cu cte ca( select order_id, customer_id, order_dt, lag(order_dt) over(partition by customer_id order by order_dt) ca prev_order_ft, datediff(day, prev_order_ft, order_dt) ca gap_days din Orders ) selectați customer_id, max(gap_days) ca gap_days din cte unde gap_days > 60 grupuți după customer_id
Ce tip de tipizare este folosit în Python: statică sau dinamică?
Cum să organizezi citirea aceleiași mesaje din Kafka de către mai mulți consumatori diferiți (fan-out)?
În ce echipă și sub ce conducere te simți cel mai confortabil să lucrezi?
Care a fost rolul tău într-o echipă de 4 persoane și cum îți vezi cariera peste 3-4 ani, ce obiective îți propui?
Cum se încarcă date din Kafka în ClickHouse prin streaming pentru rapoarte în timp real?
Povestiți despre experiența dvs. de migrare de la Spring Boot 2 la Spring Boot 3
Cum se scrie un JSON într-un folder (în contextul Airflow DAG)
Cum creezi mai multe sarcini identice în paralel în Airflow (de exemplu, încărcarea multor fișiere identice)?
În ce cazuri se aplică normalizarea și în care denormalizarea?
Ce tehnologii ai folosit pentru a obține și încărca datele în Parquet și ce mecanisme au fost utilizate?
Ce este versionarea semantică? Când să crești major, minor, patch?
Ai lucrat cu funcții de fereastră? Ce tipuri de funcții de fereastră cunoști?