升级Apache Beam 2.20.0后Airflow无法获取Dataflow作业状态问题
我来帮你拆解这个问题——这个URL错误其实是Apache Beam 2.20.0版本和Airflow的Dataflow集成之间的兼容性问题导致的,具体成因可以分成这几点来看:
Beam 2.20.0的API资源标识格式变更
Beam 2.20.0开始对Dataflow作业的资源命名格式做了调整,区域化作业的资源名称从之前的projects/{project}/jobs/{job_id},变成了带区域的完整路径projects/{project}/locations/{location}/jobs/{job_id}。这个变更本身是符合Google Cloud API规范的,但给旧版本的Airflow集成带来了适配问题。旧版Airflow的参数解析逻辑滞后
如果你使用的是Airflow 1.x系列较旧的版本(比如1.10.x早期版本),它的DataFlowHook在解析Beam返回的作业元数据时,没有适配新的资源名称格式:- 原本应该从资源名称的最后一段提取
job_id - 但旧的解析逻辑错误地截取了倒数第二段的区域值(也就是
us-central1),导致把区域当成了作业ID拼接到API URL里,就出现了你看到的/jobs/us-central1这种错误路径。
- 原本应该从资源名称的最后一段提取
作业状态查询的字段映射错误
还有一种可能是,Beam 2.20.0在返回作业提交结果时,调整了响应中jobId字段的位置或者嵌套结构,而旧版Airflow的钩子代码没有同步更新,错误地读取了location字段的值来填充作业ID参数,最终生成了错误的GET请求URL。
举个更直观的例子:
正确的作业资源名称:
projects/umg-de/locations/us-central1/jobs/abc123xyz
旧Airflow解析时错误截取了us-central1作为job_id,而不是最后的abc123xyz
本质上来说,这就是Beam升级后API响应结构和资源标识格式发生了变化,但对应的Airflow Dataflow组件没有及时跟进适配,导致参数提取错误,进而触发了这个奇怪的URL问题。
内容的提问来源于stack exchange,提问作者tank

