使用Apache Beam将CSV转换为按房产分组的换行分隔JSON
解决方案:用Apache Beam实现房产交易数据分组并生成换行JSON
错误原因
你遇到的NotImplementedError: 'ngroup' is not implemented yet是因为Apache Beam的DataFrame API是Pandas API的模拟实现,但并非所有Pandas方法都已支持,ngroup就是暂未实现的功能。需要改用Apache Beam原生的转换操作完成需求。
实现步骤
- 构造房产唯一标识:用
postcode和house_number组合成唯一键(拼接字符串或用元组),确保同一房产的交易对应同一个标识 - 读取并解析CSV数据为字典格式的交易记录
- 将每条记录转换为
(房产键, 交易记录)的键值对 - 用
GroupByKey按房产键分组,得到每个房产的所有交易列表 - 将分组结果转换为目标JSON结构:每个对象包含
property_id(唯一键)和transactions(交易数组) - 写入换行分隔的JSON文件
完整代码示例
import apache_beam as beam import csv import json from typing import Dict, Tuple def parse_csv_row(row: str) -> Dict: """解析CSV行到字典格式""" reader = csv.DictReader([row]) return next(reader) def create_property_key(transaction: Dict) -> Tuple[str, Dict]: """生成房产唯一键,返回(键, 交易记录)格式""" # 用邮编+门牌号拼接成唯一ID,可根据实际调整分隔符 property_id = f"{transaction['postcode']}_{transaction['house_number']}" return (property_id, transaction) def format_to_property_json(grouped: Tuple[str, iter]) -> str: """将分组结果转换为目标JSON字符串""" property_id, transactions = grouped property_obj = { "property_id": property_id, "transactions": list(transactions) } return json.dumps(property_obj) if __name__ == "__main__": with beam.Pipeline() as p: transactions = ( p | "读取CSV文件" >> beam.io.ReadFromText("example.csv", skip_header_lines=1) | "解析CSV行" >> beam.Map(parse_csv_row) | "生成房产键" >> beam.Map(create_property_key) | "按房产分组" >> beam.GroupByKey() | "格式化为JSON" >> beam.Map(format_to_property_json) | "写入JSON文件" >> beam.io.WriteToText( "output/property_transactions", file_name_suffix=".json", shard_name_template="" # 可选:避免生成分片文件,只输出一个文件 ) )
代码说明
- parse_csv_row:将每行CSV字符串解析为字典,保留所有交易字段。如果CSV格式复杂,建议改用
beam.io.ParseCSV指定字段类型和表头,性能更优。 - create_property_key:通过
postcode和house_number生成唯一房产ID,确保同一房产的交易被分到同一组。 - format_to_property_json:将分组后的键值对转换为要求的JSON结构,并用
json.dumps序列化为字符串,便于写入文件。 - WriteToText:生成换行分隔的JSON文件,每行对应一个房产的JSON对象。
可选优化
- 替换解析逻辑为
beam.io.ParseCSV(适合大文件):from apache_beam.io.csv import ParseCSV # 替换"解析CSV行"步骤 | "解析CSV" >> ParseCSV( skip_header_lines=1, schema="postcode:str, house_number:str, transaction_date:str, price:float" # 按实际字段定义 ) - 在
create_property_key中添加空值校验,避免生成无效的房产ID。
内容的提问来源于stack exchange,提问作者sk97
相关产品推荐
相关产品推荐

