Sobes.tech

Data Engineer

User_id боюнча skew (шикелүү) учурда транзакциялар жана колдонуучулар таблицаларын JOIN кылганда кандай иш алып баруу керек? Оптимал бөлүштүрүү ачкычын кантип тандап алса болот?

Middle
210

clickhouse_driver'nu import eting 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
197

сол жактан кошуу: 10 жазу менен таблица менен 100 жазу менен таблица. Эмне минимал жана максимал саптар саны алынышы мүмкүн?

Middle
193

Бөлүү жана таратуу (sharding) эмне?

Middle
186

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

Middle
175

from 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
174

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

Nested Loop Join учурунда кайсы ресурс эң көп сарпталат?

Middle
154

raw.local_transactions таблицасын cluster_4x2 топологиясында түзү ( transaction_id String, user_id String, amount Float64, created_at DateTime ) Мотор = ReplicatedMergeTree() Бөлүмдөргө бөлүү amount боюнча Order боюнча (transaction_id); raw.transactions таблицасын cluster_4x2 топологиясында түзү ( transaction_id String, user_id String, amount Float64, created_at DateTime ) Мотор = 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
149

ROWS BETWEEN жана RANGE BETWEEN ортосундагы айырма эмнеде?

Middle
146

Терезе функциясы SUM() OVER (PARTITION BY user_id) — бир учурда сатып алуу датасын ORDER BY кошуп жатабыз, экинчи учурда — кошпойбуз. Эмне айырмасы?

Middle
138

Маалыматтар базаларында кайсы физикалык JOIN алгоритмдер иштейт?

Middle
132

NULL plus 5 — канча болот?

Middle
131