Sobes.tech

Data Engineer

Wie vorgehen bei Skew (Schieflage) bezüglich user_id beim JOIN von Transaktions- und Benutzertabellen? Wie wählt man den optimalen Verteilungsschlüssel?

Middle
208

importiere 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: Tabelle mit 10 Einträgen LEFT JOIN Tabelle mit 100 Einträgen. Was ist die minimale und maximale Anzahl an Zeilen, die man erhalten kann?

Middle
189

Was ist Partitionierung und Verteilung (Sharding)?

Middle
186

CREATE TABLE core.localUserMetadata ON CLUSTER cluster_4x2 ( UserId String, Country String, LastLoginDate DateTime ) ENGINE = ReplikatedReplacingMergeTree() 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
173

aus 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
172

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

Welche Ressource wird bei einem Nested Loop Join am meisten verbraucht?

Middle
152

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

ERSTELLEN SIE TABELLE raw.local_transactions AUF DEM CLUSTER cluster_4x2 ( transaction_id String, user_id String, amount Float64, created_at DateTime ) ENGINE = ReplicatedMergeTree() PARTITION BY amount ORDER BY (transaction_id); ERSTELLEN SIE TABELLE raw.transactions AUF DEM CLUSTER cluster_4x2 ( transaction_id String, user_id String, amount Float64, created_at DateTime ) ENGINE = Distributed('cluster_4x2', 'raw', 'local_transactions', cityHash64(user_id));

Middle
148

Was ist der Unterschied zwischen ROWS BETWEEN und RANGE BETWEEN?

Middle
145

Die Fensterfunktion SUM() OVER (PARTITION BY user_id) — in einem Fall fügen wir ORDER BY Kaufdatum hinzu, im anderen nicht. Was ist der Unterschied?

Middle
136

Welche physischen JOIN-Algorithmen funktionieren in Datenbanken?

Middle
132

NULL plus 5 — wie viel wird das sein?

Middle
130