Data Engineer
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
Funcția de fereastră SUM() OVER (PARTITION BY user_id) — într-o situație adăugăm ORDER BY data achiziției, în alta nu. Care este diferența?
În ce mediu a fost lansat Spark?
Ce tip de stocare are Parquet: pe rând sau pe coloană?
Care sunt dezavantajele UDF-urilor în Spark și care este cel mai mare dezavantaj al lor?
Ce te atrage în domeniul travel? Poate îți place să călătorești tu însuți, sau travel-tech te-a interesat într-un fel sau altul?
Vorbește despre tipurile de date mutabile și imutabile în Python și oferă exemple.
Pe ce câmpuri a avut loc unirea datelor și cum au fost eliminate duplicatele din vitrină?
Ai lucrat cu încărcare în lot sau streaming?
De ce nu se poate folosi o listă ca cheie a unui dicționar?
Ce imagini Docker a fost nevoie să construiți și de ce?
Pentru fiecare rând din tabelul salaries, afișați numele angajatului, data, salariul și, folosind o funcție de fereastră, adăugați o coloană "salariul mediu pe departament la această dată".
Cum sunt legate sortarea din interiorul funcției de fereastră și sortarea finală ORDER BY din interogare?
print(len(' '.join(map(str, [0, 1]))))
Cum să gestionezi resursele unei aplicații Spark pentru a nu depăși memoria atunci când modifici volumul datelor de intrare?
Ce este partiționarea și distribuția (sharding)?
dict1 = { (1, 2), [3, 4, 5] : 0 } var = 1, 2
Problema #2 Scrieți o interogare care afișează TOP-3 angajați după salariu în fiecare departament. Afișați numele departamentului, numele angajatului și salariul său. tabelul employee: | id | name | salary | dept_id | |----|-------|--------|---------| | 10 | John | 1000 | 2 | | 20 | Tom | 5000 | 3 | | 30 | Bill | 3000 | 2 | tabelul department: | id | name | |----|-----------| | 1 | Marketing | | 2 | IT | | 3 | Finance | Ieșire: department_name, num, employee_name, salary cu cte ca( select d.name as department_name, e.name as employee_name, e.salary, dense_rank from employee e join department d on e.id = d.id
Povestește despre o sarcină interesantă din companie în ultimele 2,5 ani, de exemplu legată de ClickHouse.
Dă un exemplu în care ai folosit inteligența artificială pentru a automatiza procesele repetitive.