Data Engineer
Реци о индексима у Greenplum и ClickHouse
Veliko skladištenje — koliko je bilo teško održavati ga ručno u volumenima i objektima?
Ispričajte svoje iskustvo sa Airflow (orkestratorom). Čime ste koristili?
Šta je Adaptivno izvršavanje upita (AQE)?
Šta je CTE (Zajednička Izraz Tabele)?
Koliko su realistični planovi kompanije za izgradnju sistema za prikupljanje, normalizaciju i analizu podataka?
Zadatak #3 Potražite klijente kod kojih je razmak između narudžbi veći od 60 dana. Za te klijente prikazati maksimalni razmak (gap_days). Tabela narudžbi: | order_id | customer_id | order_dt | |----------|-------------|------------| | 1 | 1 | [phone] | | 2 | 1 | [phone] | | 3 | 1 | [phone] | | 4 | 1 | [phone] | | 5 | 2 | [phone] | | 6 | 2 | [phone] | | 7 | 3 | [phone] | Konačni rezultat: | customer_id | gap_days | |-------------|----------| | 1 | 123 | | 2 | 120 |
Да ли су постојали случајеви интеракције са релационим базама података?
Šta je CROSS JOIN? Radio si s tim?
Koja je razlika između == i is u Pythonu?
Zašto sada razmatraš nove ponude?
Šta je sledljivost? Imate li iskustva sa tim?
Kako razdeliti zahtev na faze? Šta za to radimo?
Na koju platu računaš?
Koje su vrste JOIN u SQL (logičke i fizičke)?
Od čega se obično sastoje logovi zadataka Airflow?
Kako dodati pomični prosek sume narudžbina za tekući i prethodna dva dana po korisniku?
Šta je shuffle u Spark-u i čemu služi?
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
Kako swap fajlovi utiču na performanse? Kada još mogu nastati swap fajlovi?