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

Apache Beam Python订单统计求助:统计金额分档的数量与总额

问题分析与修正方案

原代码的核心问题

你当前的代码在处理每一行数据时,都会新建独立的列表存储当前行的金额,然后计算这单个数据的"数量"(总是1或0)和"总和"(即当前金额本身),最后返回这四个值。这种方式只能输出每行的单独数据,无法实现全局汇总统计——因为Beam是分布式处理框架,每行数据单独处理,必须通过分组聚合才能得到全局的统计结果。

修正后的代码

import apache_beam as beam
from apache_beam.io import ReadFromText

from google.colab import drive
drive.mount('/content/drive')

orders = '/content/drive/MyDrive/ColabNotebooks/whatever/taxiorders.csv'

# 给每行金额标记分类,输出(类别,金额)键值对
class ClassifyAmount(beam.DoFn):
    def process(self, element):
        # 提取最后一列并转为浮点数,处理可能的异常
        try:
            amount = float(element.split(",")[-1])
            if amount < 15:
                yield ('小于15', amount)
            else:
                yield ('大于等于15', amount)
        except ValueError:
            # 跳过转换失败的行(比如非数字内容)
            pass

# 自定义合并函数,同时统计每个类别的数量和金额总和
class CountSumCombineFn(beam.CombineFn):
    def create_accumulator(self):
        return (0, 0.0)  # 初始化:(计数, 总和)
    
    def add_input(self, accumulator, input):
        count, total = accumulator
        return (count + 1, total + input)
    
    def merge_accumulators(self, accumulators):
        counts, totals = zip(*accumulators)
        return (sum(counts), sum(totals))
    
    def extract_output(self, accumulator):
        return {'数量': accumulator[0], '金额总和': round(accumulator[1], 2)}

with beam.Pipeline() as p:
    (p
     | '读取CSV文件' >> ReadFromText(orders, skip_header_lines=2)
     | '分类金额' >> beam.ParDo(ClassifyAmount())
     | '全局统计数量与总和' >> beam.CombinePerKey(CountSumCombineFn())
     | '转换为目标格式' >> beam.Map(lambda x: {x[0]: x[1]})
     | '打印结果' >> beam.Map(print))

代码说明

  1. ClassifyAmount 类:负责提取每行最后一列的金额并转为浮点数,给金额打上分类标签,输出(分类标签, 金额)的键值对,为后续分组聚合做准备。
  2. CountSumCombineFn 类:自定义Beam合并函数,对每个分类下的所有金额数据,同时统计数量和总和,最终返回包含这两个指标的字典。
  3. Pipeline流程:读取文件→分类标记→按分类聚合统计→转换为目标键值对格式→输出结果。

额外提示

如果你的CSV文件存在字段包含逗号的情况(比如带引号的文本字段),直接用split(",")会导致拆分错误,建议改用Python内置的csv模块解析,示例代码如下:

import csv
from io import StringIO

# 在ClassifyAmount的process方法中替换拆分逻辑
def process(self, element):
    reader = csv.reader(StringIO(element))
    row = next(reader)
    amount = float(row[-1])
    # 后续分类逻辑不变

内容的提问来源于stack exchange,提问作者Pancho

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 19:16:26