Data Engineer
Wie man einen gleitenden Durchschnitt der Bestellsumme des aktuellen und der beiden vorherigen Tage pro Benutzer hinzufügt?
Wie funktioniert der Hash-Join für Tabellen mit 1.000 und 1.000.000 Zeilen?
Wozu brauchten Sie überhaupt Data Vault?
Wurde die Methode des Vergleichs von Hashes von Datensätzen verwendet, um die Differenz zwischen Quelle und Ziel Tabelle zu berechnen?
Haben Sie beim Einfügen den Fehler ClickHouse `Too many parts` erlebt? Wie haben Sie ihn gelöst?
Mit welcher Version von Airflow haben Sie gearbeitet?
Mit welchen Dateiformaten haben Sie gearbeitet?
Welche Arten von physischen JOINs gibt es in Spark?
Was passiert physisch bei INSERT, UPDATE und DELETE in Greenplum; warum wird nach DELETE der Festplattenspeicher nicht freigegeben und worin unterscheidet sich DELETE von TRUNCATE?
Warum können Daten verloren gehen? Nennen Sie einige Gründe.
Warum ziehst du jetzt neue Angebote in Betracht?
Erzählen Sie mir von sich: Erfahrungen, Aufgaben, Funktionen oder Errungenschaften, auf die Sie stolz sind.
Haben Sie Erfahrung mit FastAPI?
Wie haben Sie bei der erneuten Ausführung des DAG nach einem Absturz mit den Dateien gearbeitet, um keine unvollständigen Daten in S3 zu überschreiben?
In welcher Situation kann die Suche in einem Wörterbuch auf den schlechtesten Fall verschlechtern?
Wie haben Sie Optimierungsprobleme gelöst, wenn eine Abfrage lange dauert und einen Vollscan durchführt? Wie würden Sie an eine neue solche Aufgabe herangehen?
Was ist Datenpartitionierung und wie hilft sie, eine vollständige Tabellensuche zu vermeiden?
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
Was bewirkt der Parameter depends_on_past in Airflow?
Was ist der Unterschied zwischen ACID und dem CAP-Theorem?