Apache Beam使用BigQuerySource未生成datasetName对应displayData问题咨询
问题原因
你当前使用的是BigQuerySource的查询模式(传入了query参数而非直接指定table参数),Apache Beam原生的BigQuerySource的display_data实现逻辑中,只有显式传入dataset.table格式的表路径作为table参数时,才会自动解析出datasetName、tableName这类字段并加入displayData对象。自定义查询场景下,原生实现不会解析SQL语句中的表信息,因此不会生成对应的datasetName字段。
解决方法
你可以根据自身使用场景选择以下任意一种方案实现需求:
- 方案1:改用指定表路径的读取模式
如果你的读取逻辑不需要自定义SQL,直接指定表路径构造BigQuerySource即可让框架自动生成符合要求的displayData:
bigqcollection = p | 'ReadFromBQ' >> beam.io.Read( beam.io.BigQuerySource( project=project, table="你的项目ID:你的数据集名.你的表名" ) )
- 方案2:自定义PTransform封装读取逻辑手动添加displayData
如果你必须使用自定义查询模式,可以自己实现一个包装PTransform,重写display_data方法添加你需要的字段:
class CustomBQRead(beam.PTransform): def __init__(self, project, query, dataset_name, table_name=None): self.project = project self.query = query self.dataset_name = dataset_name self.table_name = table_name def expand(self, pcoll): return pcoll | beam.io.Read(beam.io.BigQuerySource( project=self.project, query=self.query )) def display_data(self): res = { "datasetName": self.dataset_name } if self.table_name: res["table"] = self.table_name return res # 调用示例 bigqcollection = p | 'ReadFromBQ' >> CustomBQRead( project=project, query=get_java_query, dataset_name="你需要展示的数据集名", table_name="你需要展示的表名" # 可选 )
上述自定义PTransform返回的displayData会完全匹配你期望的结构,namespace字段会自动填充为该自定义类的路径。
内容的提问来源于stack exchange,提问作者Mantas Savaniakas
相关产品推荐
相关产品推荐

