Data Engineer
Vertel eens over jezelf: ervaring, taken, functies of prestaties waar je trots op bent.
Heb je ervaring met FastAPI?
Hoe lang hebt u in totaal met Spark gewerkt en hoe diep bent u erin gedoken?
dict1 = { (1, 2), [3, 4, 5] : 0 } var = 1, 2
Vertel over een interessante taak in het bedrijf in de afgelopen 2,5 jaar, bijvoorbeeld gerelateerd aan ClickHouse.
In welke situatie kan zoeken in een woordenboek degraderen tot het slechtste geval?
Hoe hebt u optimalisatieproblemen opgelost wanneer een query lang duurt en een volledige scan uitvoert? Hoe zou u een nieuwe dergelijke taak aanpakken?
Hoe voeg je een voortschrijdend gemiddelde toe van de som van bestellingen voor de huidige en de twee voorgaande dagen per gebruiker?
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
Wat doet de parameter depends_on_past in Airflow?
Waarvoor had je in het algemeen Data Vault nodig?
Vertel over de replicatie-instelling die je hebt opgelost.
Is de methode van het vergelijken van hash-registraties gebruikt om het verschil tussen de bron en de doel tabel te berekenen?
Wat gebeurde er toen de datakwaliteitsbeperking werd geactiveerd: waarschuwing, crash, herstart?
Vergelijk de ETL- en ELT-aanpakken: wat is het verschil en wanneer wordt elk van hen toegepast?
Hoe werkten ze met 1C: haalden ze gegevens rechtstreeks uit de database, via een bus of op een andere manier?
Met welke databases hebt u gewerkt? Hoe werkt MongoDB en hoe schaalt het?
Hoe hebben jullie vastgesteld dat de invoer al bestaat en niet opnieuw hoeft te worden opgeslagen?
NULL plus 5 — hoeveel wordt dat?
Kan je in Spark parallel gegevens lezen uit één Parquet-bestand, of alleen met één kern in een taak?