Sobes.tech

Data Engineer

Cum să procedezi în cazul unui skew (părtinire) pe user_id la JOIN între tabelele de tranzacții și utilizatori? Cum alegi cheia de distribuție optimă?

Middle
213

importă 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
198

LEFT JOIN: tabel cu 10 înregistrări LEFT JOIN tabel cu 100 înregistrări. Care este numărul minim și maxim de rânduri pe care le poți obține?

Middle
194

Ce este partiționarea și distribuția (sharding)?

Middle
186

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

Middle
178

din 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
177

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
161

Care resursă este consumată cel mai mult la Nested Loop Join?

Middle
154

CREARE TABELA raw.local_transactions PE CLUSTERUL cluster_4x2 ( transaction_id String, user_id String, amount Float64, created_at DateTime ) MOTOR = ReplicatedMergeTree() PARTITIONARE DUPĂ amount ORDONARE DUPĂ (transaction_id); CREARE TABELA raw.transactions PE CLUSTERUL cluster_4x2 ( transaction_id String, user_id String, amount Float64, created_at DateTime ) MOTOR = Distributed('cluster_4x2', 'raw', 'local_transactions', cityHash64(user_id));

Middle
151

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
150

Care este diferența dintre ROWS BETWEEN și RANGE BETWEEN?

Middle
147

Funcția de fereastră SUM() OVER (PARTITION BY user_id) — într-o situație adăugăm ORDER BY data achiziției, în alta nu. Care este diferența?

Middle
138

Ce algoritmi fizici JOIN funcționează în bazele de date?

Middle
133

NULL plus 5 — cât va fi?

Middle
131