【发布时间】:2023-02-14 16:16:43
【问题描述】:
我有一个气流任务来用这个流提取数据
PostgreSQL -> Google Cloud Storage -> BigQuery
我遇到的问题是,似乎并非所有数据都被提取到 BigQuery 中。在 PostgreSQL 源上,该表有 18M+ 行数据,但在摄取后它只有 4M+ 行数据。
当我检查生产时,数据通过此查询返回 18M+ 行:
SELECT COUNT(1) FROM my_table
-- This return 18M+ rows
但是在 DAG 完成运行后,当我检查 BigQuery 时:
SELECT COUNT(1) FROM data_lake.my_table
-- This return 4M+ rows
为了做笔记,并非我摄取的所有表格都是这样返回的。所有较小的桌子都摄取得很好。但是当它达到一定数量的行时,它的行为就像这样。
我怀疑数据是从 PostgreSQL 提取到 Google Cloud Storage 的。所以我会在这里提供我的功能:
def create_operator_write_append_init(self, worker=10):
worker_var = dict()
with TaskGroup(group_id=self.task_id_init) as tg1:
for i in range(worker):
worker_var[f'worker_{i}'] = PostgresToGCSOperator(
task_id = f'worker_{i}',
postgres_conn_id = self.conn_id,
sql = 'extract_init.sql',
bucket = self.bucket,
filename = f'{self.filename_init}_{i}.{self.export_format}',
export_format = self.export_format, # the export format is json
gzip = True,
params = {
'worker': i
}
)
return tg1
这是 SQL 文件:
SELECT id,
name,
created_at,
updated_at,
deleted_at
FROM my_table
WHERE 1=1
AND ABS(MOD(hashtext(id::TEXT), 10)) = {{params.worker}};
我所做的是将数据分块并将其分成几个工作人员,因此是 TaskGroup。
提供更多信息。我使用作曲家:
-
作曲家-2.0.32-气流-2.3.4
-
大型实例
-
工人 8CPU
-
工人 32GB 内存
-
工人 2GB 存储空间
-
1-16 岁之间的工人
这些发生的可能性有多大?
【问题讨论】:
标签: python airflow google-cloud-composer airflow-2.x