Data Engineer
Jak rozwiązywałeś problemy optymalizacyjne, gdy zapytanie działa długo i wykonuje pełny skan? Jak podszedłbyś do nowego takiego zadania?
Jak dodać ruchomą średnią sumy zamówień za bieżący i dwa poprzednie dni na użytkownika?
Potrzebujesz szybko przywrócić kilka plików do wersji z ostatniego commita, nie naruszając innych zmian. Jak postępować? git reset --hard HEAD Usuń i utwórz pliki ręcznie git fetch i git merge git revert HEAD git checkout HEAD^ <plik1> <plik2>
Jakie zadanie w Pythonie rozwiązałeś ostatnio w produkcji?
Jak rozwiązać problem, gdy proces odczytu może zobaczyć pośredni (niezgodny) stan tabeli podczas wieloetapowego ETL (delete+insert) w Iceberg?
Na wyświetlonym ekranie znajdują się dwie struktury danych; w której z nich wyszukiwanie wartości 2 będzie szybsze i dlaczego?
Podaj przykład typowego zadania związane z pisaniem DAG w Airflow.
Co się stało i jak odzyskać pracę?
Dlaczego oddziela się widoki (VIEW) i tabele? Do czego służą widoki?
Czym jest ACID? Co oznacza każda litera? Jak to się odnosi do transakcji?
Jak napisałeś DAG Airflow i zadania — ręcznie czy za pomocą szablonów?
Czym różni się EXPLAIN ANALYZE od zwykłego EXPLAIN? Dlaczego nie zaleca się uruchamiania go na zapytaniach DELETE/UPDATE?
Czym jest semantyczne wersjonowanie? Kiedy zwiększać major, minor, patch?
Czy możesz opowiedzieć, co ostatecznie tutaj uzyskamy? Należy skomentować każdy krok, jak to działa.
Jak Twoje doświadczenie odpowiada obowiązkom na stanowisku: wsparcie i analiza procesów ładowania, monitorowanie i wykrywanie anomalii, testowanie i wdrażanie poprawek w produkcji, wsparcie drugiej linii, sprawdzanie dokumentacji technicznej?
Na jakim etapie, jakie kontrole są przeprowadzane i co się dzieje z nieprawidłowymi danymi podczas kontroli jakości danych (Data Quality)?
Z czego zazwyczaj składają się logi zadań Airflow?
Czym jest plan zapytań? Do czego są potrzebne?
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
Jak są powiązane partycje Spark i partycje w tabelach (np. w Iceberg)?