Azure Dataflow中如何将列转换为另一列的JSON对象
使用Dataflow将记录格式化为JSON并写入新列
实现逻辑:通过Dataflow的自定义转换,把每行的指定字段(或全部字段)转换成JSON字符串,再作为新列附加到原记录中。
Python SDK示例代码:
- 导入依赖库
import json import apache_beam as beam from apache_beam.transforms import DoFn - 定义转换处理类
class AddJsonColumn(DoFn): def process(self, row): # 提取需要转JSON的字段(排除目标新列,若要全字段可直接用dict(row)) json_content = {key: row[key] for key in row if key != "formatted_json"} # 转成JSON字符串并添加为新列 row["formatted_json"] = json.dumps(json_content) yield row - 构建Dataflow管道
with beam.Pipeline() as pipeline: ( pipeline # 读取输入数据(替换为你的数据源,比如BigQuery、CSV) | "读取输入" >> beam.io.ReadFromBigQuery(query="SELECT * FROM `project.dataset.input_table`") # 执行JSON转换并添加新列 | "生成JSON列" >> beam.ParDo(AddJsonColumn()) # 写入输出(替换为你的输出目标) | "写入输出表" >> beam.io.WriteToBigQuery( table="project.dataset.output_table", write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) )
- 导入依赖库
关键说明:
- 若输入是CSV等非结构化文本,需先通过
beam.Map将每行解析为字典/Row对象。 - 如果只需要特定字段生成JSON,修改
json_content的字段筛选规则即可。 - 输出到BigQuery时,确保表结构包含
formatted_json字符串类型的列,或开启自动.schema创建。
- 若输入是CSV等非结构化文本,需先通过
内容的提问来源于stack exchange,提问作者Pankaj Prabhakar
相关产品推荐
相关产品推荐

