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

如何在Apache Beam中用Combine.perKey实现设备与类型的复合键聚合

修改方案:生成device + measurementType复合键

要实现按设备名称和测量类型的复合键聚合,你需要把KV的键从单一的measurementType字符串,改成包含device和measurementType的复合结构。下面提供几种可行的实现方式:

方式1:自定义不可变复合键类(推荐)

创建一个包含device和measurementType的不可变类,必须重写equals()和hashCode()方法——这是Beam中KV键的硬性要求,确保能正确按组合字段分组:

import java.io.Serializable;
import java.util.Objects;

public class DeviceMeasurementKey implements Serializable {
    private final String device;
    private final String measurementType;

    public DeviceMeasurementKey(String device, String measurementType) {
        this.device = device;
        this.measurementType = measurementType;
    }

    // Getter方法
    public String getDevice() { return device; }
    public String getMeasurementType() { return measurementType; }

    // 重写equals和hashCode保证分组正确性
    @Override
    public boolean equals(Object o) {
        if (this == o) return true;
        if (o == null || getClass() != o.getClass()) return false;
        DeviceMeasurementKey that = (DeviceMeasurementKey) o;
        return Objects.equals(device, that.device) && Objects.equals(measurementType, that.measurementType);
    }

    @Override
    public int hashCode() {
        return Objects.hash(device, measurementType);
    }
}

然后修改你的DoFn实现,输出的KV键替换为这个自定义类:

public class KvByDeviceAndMeasurementType extends DoFn<Measurement, KV<DeviceMeasurementKey, Measurement>> implements Serializable {
    @ProcessElement
    public void processElement(ProcessContext context) {
        Measurement measurement = context.element();
        DeviceMeasurementKey compositeKey = new DeviceMeasurementKey(
            measurement.getDevice(),
            measurement.getMeasurementType()
        );
        context.output(KV.of(compositeKey, measurement));
    }
}

后续使用Combine.perKey()时,Beam会自动按device和measurementType的组合分组计算平均值。

方式2:使用第三方Tuple类型(快速实现)

如果不想自定义类,可以使用Apache Commons Lang的Pair类(需引入commons-lang3依赖),它已经实现了序列化、equals和hashCode:

import org.apache.commons.lang3.tuple.Pair;

public class KvByDeviceAndMeasurementType extends DoFn<Measurement, KV<Pair<String, String>, Measurement>> implements Serializable {
    @ProcessElement
    public void processElement(ProcessContext context) {
        Measurement measurement = context.element();
        Pair<String, String> compositeKey = Pair.of(measurement.getDevice(), measurement.getMeasurementType());
        context.output(KV.of(compositeKey, measurement));
    }
}

方式3:拼接字符串(不推荐)

如果只是临时测试,可以把两个字段拼接成单一字符串当键,比如用特殊分隔符区分:

public class KvByDeviceAndMeasurementType extends DoFn<Measurement, KV<String, Measurement>> implements Serializable {
    @ProcessElement
    public void processElement(ProcessContext context) {
        Measurement measurement = context.element();
        // 注意:如果字段本身包含分隔符(比如|),会导致分组错误
        String compositeKey = measurement.getDevice() + "|" + measurement.getMeasurementType();
        context.output(KV.of(compositeKey, measurement));
    }
}

后续聚合说明

修改完DoFn后,Combine.perKey()的逻辑无需大幅调整,只需将键的类型对应更新即可,比如使用自定义键类时,聚合函数的输入键类型为DeviceMeasurementKey。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 12:33:19