Data Engineer
Parlez-moi de vous : expérience, tâches, fonctionnalités ou réalisations dont vous êtes fier.
Avez-vous de l'expérience avec FastAPI?
Combien de temps avez-vous travaillé avec Spark au total et à quel point vous y êtes-vous impliqué?
dict1 = { (1, 2), [3, 4, 5] : 0 } var = 1, 2
Parlez d'une tâche intéressante dans l'entreprise au cours des 2,5 dernières années, par exemple liée à ClickHouse.
Dans quelle situation la recherche dans un dictionnaire peut-elle se dégrader jusqu'au pire cas?
Comment avez-vous résolu les problèmes d'optimisation lorsque la requête prend beaucoup de temps et effectue une analyse complète ? Comment aborderiez-vous une nouvelle tâche similaire ?
Comment ajouter une moyenne mobile de la somme des commandes des deux jours précédents et du jour actuel par utilisateur?
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
Que fait le paramètre depends_on_past dans Airflow?
À quoi vous servait en général Data Vault?
Parlez-moi de la tâche de configuration de la réplication que vous avez résolue.
La méthode de comparaison des hachages des enregistrements a-t-elle été utilisée pour calculer la différence entre la source et la table cible?
Que se passait-il lorsque la restriction de qualité des données était déclenchée : alerte, panne, redémarrage ?
Comparez les approches ETL et ELT : quelle est la différence et quand chacune d'elles est-elle utilisée ?
Comment travaillaient-ils avec 1C : récupéraient-ils les données directement de la base de données, via un bus ou d'une autre manière ?
Avec quelles bases de données avez-vous travaillé ? Comment fonctionne MongoDB et comment se scale-t-elle ?
Comment avez-vous déterminé que l'enregistrement existe déjà et qu'il n'est pas nécessaire de le réenregistrer?
NULL plus 5, combien cela fera-t-il?
Peut-on lire des données en parallèle à partir d'un seul fichier Parquet dans Spark, ou seulement avec un seul noyau dans une tâche?