You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 01:47:42