Airflow任务执行成功但数据未全量同步至BigQuery的排查问询
可能的原因分析
1. 分片SQL的哈希逻辑存在数据遗漏
- NULL值处理缺失:如果源表
my_table中id字段存在NULL值,hashtext(id::TEXT)会返回NULL,导致WHERE条件不匹配,这部分数据不会被任何分片抽取。可执行以下查询验证:SELECT COUNT(1) FROM my_table WHERE id IS NULL; - 哈希分布异常:虽然
ABS(MOD(hashtext(id::TEXT), 10))理论上会将数据分配到0-9的分片,但hashtext的哈希结果可能存在极端分布,比如某几个分片无数据。可执行以下查询检查分片数据分布:SELECT ABS(MOD(hashtext(id::TEXT), 10)) AS shard, COUNT(1) FROM my_table GROUP BY shard ORDER BY shard;
2. 部分分片任务未正常执行
- Worker存储资源不足:Worker仅配置2GB存储,抽取大分片时生成的gzip临时文件可能超出存储上限,导致任务崩溃但Airflow未标记为失败。需检查每个
worker_i任务的日志,确认是否存在存储溢出情况。 - 任务并发限制:Airflow的
dag_concurrency或worker_concurrency配置可能限制了同时执行的分片任务数量,导致部分任务排队超时被跳过。需查看TaskGroup下所有任务的状态,确认10个分片任务是否全部成功执行。 - TaskGroup状态异常:个别分片任务实际失败,但TaskGroup的状态判断逻辑存在问题,导致整体显示为成功。需逐一核查每个分片任务的日志和GCS中对应文件是否存在。
3. GCS到BigQuery的加载环节遗漏文件
- 文件匹配规则错误:BigQuery加载任务可能仅匹配了部分分片文件(如仅加载
filename_init_0至filename_init_4),而非所有filename_init_*.json.gz文件。需检查BigQuery加载操作的文件路径配置。 - GCS文件权限问题:部分分片文件生成后,BigQuery服务账号无读取权限,导致加载时跳过这些文件且无报错。需验证GCS bucket的权限配置,确保BigQuery服务账号能访问所有分片文件。
- 文件未完全上传触发加载:Airflow标记分片任务成功时,GCS文件可能仍在异步上传中,导致BigQuery加载时读取到不完整文件或跳过未完成文件。需确保加载任务在所有分片任务完成后再执行。
4. PostgresToGCSOperator的内部限制或bug
- 大结果集处理异常:Airflow 2.3.4版本的
PostgresToGCSOperator在处理超大规模结果集时,可能存在内存溢出或分页截断问题,导致仅抽取部分数据。需查看任务日志中是否有内存相关的警告或错误。 - JSON序列化异常:
created_at等时间字段或特殊字符字段在序列化为JSON时出现异常,导致部分行无法写入文件,且Operator未抛出错误。可检查GCS中的分片文件,确认是否存在数据截断或格式错误。
5. 同步过程中源表数据变更
同步期间若源表发生批量删除或数据清洗操作,可能导致目标数据量减少。可对比同步前后源表的行数变化,排除该可能性。
内容的提问来源于stack exchange,提问作者Iqbal
相关产品推荐
相关产品推荐

