Data Engineer
Pentru ce sunt folosite funcțiile de fereastră?
Ce restricție există atunci când se utilizează DISTINCT ON în PostgreSQL cu ORDER BY?
Ai primit o sarcină cu numărul ABC. Care este ordinea ta de acțiuni în Git la începutul lucrului?
Ați lucrat vreodată cu decoratori în Python? Ce sunt și unde sunt folosiți?
Care este pericolul shuffle-ului în Spark?
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
Ce poți să-mi spui despre API-ul SOAP?
Care este diferența dintre ROWS BETWEEN și RANGE BETWEEN?
Cum se aplică o funcție la toate rândurile sau elementele din pandas?
La ce trebuie să fii atent în privința memoriei, vitezei și corectitudinii la parsarea fișierelor XML mari?
Cum anume ai citit Parquet?
Ai lucrat direct cu Spark? Care este dezavantajul său?
Lați proiecte Scala la care ai lucrat anterior? Împărtășește sarcini specifice și realizări în această tehnologie.
Cum ai scris DAG-ul Airflow și joburile: manual sau cu șabloane?
Raport pentru compania logistică Ești un analist al unei companii logistice care ține evidența operațiunilor în depozite. Trebuie să întocmești un raport despre eficiența fiecărui depozit. Pentru fiecare depozit, calculează: • numărul total de operațiuni (count_operations); • numărul total de produse procesate în depozit (sum_quantity); • timpul mediu de procesare a unei operațiuni (avg_processing_time), luând în considerare doar operațiunile cu timp specificat (nu NULL), rotunjit la cel mai apropiat întreg; • numărul maxim și minim de produse procesate într-o singură operațiune (max_quantity, min_quantity); • numărul de operațiuni de fiecare tip («livrare», «expediere», «transfer») în coloane separate: supply_operations, shipment_operations, transfer_operations. Filtrează depozitele al căror număr total de operațiuni este mai mare de 2 și timpul mediu de procesare nu depășește 60 de minute. Ordonează rezultatul după ID-ul depozitului în ordine crescătoare. Formatul de intrare Tabelul operations: • operation_id (int) — identificator unic al operației • warehouse_id (int) — identificatorul depozitului • operation_type (text) — tipul operației: «livrare», «expediere», «transfer»
Au fost cazuri de interacțiune cu baze de date relaționale?
Ce este versionarea semantică? Când să crești major, minor, patch?
Cu dependențe între joburi, între DAG-uri — ați folosit senzori pentru ca acestea să pornească unul după altul mai multe joburi?
Cum îți organizezi proiectele? Scrii un README sau altceva?
Când este potrivit ThreadPoolExecutor?