Sobes.tech

Data Engineer

Comment gérer le déséquilibre (skew) sur user_id lors d'un JOIN entre les tables de transactions et d'utilisateurs ? Comment choisir la clé de distribution optimale ?

Middle
208

importer clickhouse_driver from airflow.hooks.base import * def get_clickhouse_client(): conn = BaseHook.get_connection("clickhouse_default") return clickhouse_driver.Client( host=conn.host, port=conn.port, user=conn.login, password=conn.password, database=conn.schema )

Middle
195

LEFT JOIN : table avec 10 enregistrements LEFT JOIN table avec 100 enregistrements. Quel est le nombre minimum et maximum de lignes que l'on peut obtenir?

Middle
190

Qu'est-ce que le partitionnement et la distribution (sharding) ?

Middle
186

CREATE TABLE core.localUserMetadata ON CLUSTER cluster_4x2 ( UserId String, Country String, LastLoginDate DateTime ) ENGINE = ReplicatedReplacingMergeTree() ORDER BY (UserId); CREATE TABLE core.userMetadata ON CLUSTER cluster_4x2 ( UserId String, Country String, LastLoginDate Datetime ) ENGINE = Distributed('cluster_4x2', 'core', 'localUserMetadata', cityHash64(LastLoginDate));

Middle
174

depuis interview.utils import get_clickhouse_client from airflow import DAG from airflow.operators.python import PythonOperator from airflow.sensors.external_task import ExternalTaskSensor from datetime import datetime import pandas as pd import clickhouse_driver import os CLICKHOUSE_CLIENT = get_clickhouse_client() default_args = { "start_date": datetime(2024, 1, 1) } with DAG( dag_id="datamarts.daily_revenue_per_country", default_args=default_args, schedule_interval="@daily", catchup=False ) as 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( ... )

Middle
173

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

Middle
158

Quelle ressource est la plus consommée lors d'une jointure par boucle imbriquée?

Middle
153

CREATE TABLE datamarts.daily_revenue_per_country ON CLUSTER cluster_4x2 ( event_date Date, country String, total_revenue Float64 ) ENGINE = MergeTree() PARTITION BY toYYYYMM(event_date) ORDER BY (event_date);

Middle
149

CRÉER UNE TABLE raw.local_transactions SUR CLUSTER cluster_4x2 ( transaction_id String, user_id String, amount Float64, created_at DateTime ) MOTEUR = ReplicatedMergeTree() PARTITION PAR amount ORDER PAR (transaction_id); CRÉER UNE TABLE raw.transactions SUR CLUSTER cluster_4x2 ( transaction_id String, user_id String, amount Float64, created_at DateTime ) MOTEUR = Distributed('cluster_4x2', 'raw', 'local_transactions', cityHash64(user_id));

Middle
148

Quelle est la différence entre ROWS BETWEEN et RANGE BETWEEN ?

Middle
145

La fonction de fenêtre SUM() OVER (PARTITION BY user_id) — dans un cas, nous ajoutons ORDER BY date d'achat, dans l'autre non. Quelle est la différence?

Middle
136

Quels algorithmes physiques JOIN fonctionnent dans les bases de données?

Middle
132

NULL plus 5, combien cela fera-t-il?

Middle
131