Data Engineer
Jak přidat klouzavý průměr součtu objednávek za aktuální a předchozí dva dny na uživatele?
Jak funguje hash join pro tabulky s 1 000 a 1 000 000 řádky?
Na co vám vůbec byl Data Vault?
Byla použita metoda porovnání hashů záznamů k výpočtu rozdílu mezi zdrojem a cílovou tabulkou?
Setkali jste se s chybou ClickHouse `Too many parts` při vkládání? Jak jste ji řešili?
S jakou verzí Airflow jste pracoval?
S jakými formáty souborů jste pracoval?
Jaké typy fyzických JOINů existují v Spark?
Co se fyzicky děje při INSERT, UPDATE a DELETE v Greenplum; proč po DELETE není uvolněno místo na disku a v čem se DELETE liší od TRUNCATE?
Proč mohou být data ztracena? Uveďte několik důvodů.
Proč teď zvažuješ nové nabídky?
Povězte nám o sobě: zkušenosti, úkoly, funkce nebo úspěchy, na které jste hrdý.
Máte zkušenosti s FastAPI?
Jak jste pracoval se soubory při opětovném spuštění DAG po jeho selhání, aby nedošlo k přepsání neúplných dat v S3?
V jaké situaci může hledání ve slovníku degradovat na nejhorší případ?
Jak jste řešili problémy s optimalizací, když dotaz trvá dlouho a provádí úplný sken? Jak byste přistoupili k takové nové úloze?
Co je to rozdělení dat a jak pomáhá vyhnout se úplnému skenování tabulky?
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
Co dělá parametr depends_on_past v Airflow?
Jaký je rozdíl mezi ACID a CAP teoremou?