Data Engineer
def extract_from_s3(**kwargs): df = pd.read_csv("s3://mano-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
Langinio langelio funkcija SUM() OVER (PARTITION BY user_id) — vienu atveju pridedame ORDER BY pirkimo datą, kitame ne. Kuo tai skiriasi?
Kokioje aplinkoje buvo paleistas Spark?
Koks Parquet saugojimo tipas: eilutė ar stulpelis?
Kokia yra UDF Spark programoje ir koks jų minusas?
Kas tave traukia travel srityje? Galbūt pats mėgsti keliauti, arba travel-tech kažkaip tave sudomino?
Papaskinkite apie mutabilius ir nemutabilius duomenų tipus Python ir pateikite pavyzdžių.
Kuriuose laukuose įvyko duomenų sujungimas ir kaip buvo pašalinti dubliuojami elementai iš vitrinų?
Ar dirbote su batch įkėlimu ar srautiniu duomenų perdavimu?
Kodėl sąrašo negalima naudoti kaip žodyno raktą?
Kokius Docker vaizdus reikėjo sukurti ir kodėl?
Kiekvienai eilutei salaries lentelėje parodykite darbuotojo vardą, datą, atlyginimą ir naudodami lango funkciją pridėkite stulpelį "vidutinis atlyginimas pagal skyrių šią datą".
Kaip susijęs lango funkcijos viduje esantis rūšiavimas ir užklausos galutinis ORDER BY rūšiavimas?
print(len(' '.join(map(str, [0, 1]))))
Kaip valdyti Spark programos išteklius, kad neperkrautume atminties keičiant įvesties duomenų kiekį?
Kas yra skaidymas ir paskirstymas (sharding)?
dict1 = { (1, 2), [3, 4, 5] : 0 } var = 1, 2
Užduotis #2 Parašykite užklausą, kuri rodo TOP-3 darbuotojus pagal atlyginimą kiekviename skyriuje. Rodyti skyriaus pavadinimą, darbuotojo vardą ir jo atlyginimą. employee lentelė: | id | name | salary | dept_id | |----|-------|--------|---------| | 10 | John | 1000 | 2 | | 20 | Tom | 5000 | 3 | | 30 | Bill | 3000 | 2 | department lentelė: | id | name | |----|-----------| | 1 | Marketing | | 2 | IT | | 3 | Finansai | Išvestis: department_name, num, employee_name, salary su cte kaip( select d.name as department_name, e.name as employee_name, e.salary, dense_rank iš employee e jungtis su department d on e.id = d.id
Papaskinkite apie įdomią užduotį įmonėje per pastaruosius 2,5 metų, pavyzdžiui, susijusią su ClickHouse.
Papaskink savo pavyzdį, kada naudoji dirbtinį intelektą rutininėms procesams automatizuoti.