Data Engineer
Какво е Materialized View и как се различава от обикновен изглед?
Какво се случи и как да възстановите работата?
На какви проекти на Scala сте работили преди? Споделете конкретни задачи и постижения в тази технология.
Как написа DAG и задачите в Airflow — ръчно или с шаблони?
Как да изведете последното пълно име за всеки клиент, ако е известна само началната дата (крайната дата не е известна)?
Какво е семантично версиониране? Кога да увеличаваме major, minor, patch?
Кои магически методи на Python трябва да се вземат предвид?
С зависимости между работи, между DAG — използвахте сензори, за да се стартират няколко работи една след друга?
Как организирате проектите си? Пишете README или нещо друго?
Кога е подходящ ThreadPoolExecutor?
Получихте задача с номер ABC. Какъв е вашият ред на действия в Git при започване на работа?
Имал ли си ad-hoc задачи въпреки наличието на ясен технически задание, и как се отнасяш към ad-hoc заявки и промени в приоритетите?
от interview.utils импортирайте get_clickhouse_client от airflow импортирайте DAG от airflow.operators.python импортирайте PythonOperator от airflow.sensors.external_task импортирайте ExternalTaskSensor от datetime импортирайте datetime импортирайте pandas като pd импортирайте clickhouse_driver импортирайте os CLICKHOUSE_CLIENT = get_clickhouse_client() default_args = { "start_date": datetime(2024, 1, 1) } с DAG( dag_id="datamarts.daily_revenue_per_country", default_args=default_args, schedule_interval="@daily", catchup=False ) като dag: 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" ) 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( ... )
Имаме таблица с клиенти (id, име, адрес и т.н.), продукти (id, име и т.н.), движение на стоки (тук имаме id на клиента, продукт, дата, количество, сума), трябва да намерим кои продукти са продадени на 1 януари, само уникалните, и също така да покажем името на купувача и...
Как да сравним оригиналната („на живо“) и целевата таблица, ако източникът постоянно се променя?
Подробно опишете алгоритъма за улавяне на инкремента от източника, за да не се губи нищо при промяна на честотата на зареждане.
Как са свързани партициите на Spark и партициите в таблиците (например в Iceberg)?
Кога може да се използва семейството формати CSV?
Как да решим проблема, когато процесът на четене може да види междинното (несъответстващо) състояние на таблицата по време на многостепенно ETL (delete+insert) в Iceberg?
Как да приложа функция към всички редове или елементи в pandas?