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

如何在Beam Pub/Sub Pipeline中按sensorId分组固定窗口并持久化?

按SensorId分组窗口的Beam Pipeline实现方案

要实现每个sensorId对应独立固定窗口并持久化的需求,核心是将数据转换为KV结构后结合GroupByKey与窗口机制,以下是具体实现步骤:

1. 定义传感器数据模型

先创建一个POJO类映射输入的JSON格式,方便后续解析和处理:

import com.google.gson.annotations.SerializedName;

public class SensorData {
    @SerializedName("deviceId")
    private String deviceId;
    @SerializedName("value")
    private double value;
    @SerializedName("sensorId")
    private String sensorId;
    @SerializedName("timeInNanoSeconds")
    private long timeInNanoSeconds;
    @SerializedName("receivedTimeInNanoSeconds")
    private long receivedTimeInNanoSeconds;

    // Gson解析必须的无参构造函数
    public SensorData() {}

    // Getter和Setter方法
    public String getDeviceId() { return deviceId; }
    public void setDeviceId(String deviceId) { this.deviceId = deviceId; }
    public double getValue() { return value; }
    public void setValue(double value) { this.value = value; }
    public String getSensorId() { return sensorId; }
    public void setSensorId(String sensorId) { this.sensorId = sensorId; }
    public long getTimeInNanoSeconds() { return timeInNanoSeconds; }
    public void setTimeInNanoSeconds(long timeInNanoSeconds) { this.timeInNanoSeconds = timeInNanoSeconds; }
    public long getReceivedTimeInNanoSeconds() { return receivedTimeInNanoSeconds; }
    public void setReceivedTimeInNanoSeconds(long receivedTimeInNanoSeconds) { this.receivedTimeInNanoSeconds = receivedTimeInNanoSeconds; }
}

2. 解析JSON为KV结构

将Pub/Sub接收的字符串消息解析为SensorData对象,再转换为KV<String, SensorData>,其中key为sensorId——这是后续分组的依据:

import com.google.gson.Gson;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.TypeDescriptors;

// 接入原有Pipeline流程
.apply("Read PubSub Messages", PubsubIO.readStrings().fromTopic(options.getInputTopic()))
.apply("Parse JSON to KV Pair", MapElements.into(TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptors.of(SensorData.class)))
    .via(message -> {
        Gson gson = new Gson();
        SensorData data = gson.fromJson(message, SensorData.class);
        return KV.of(data.getSensorId(), data);
    }))

3. 应用窗口并按SensorId分组

先给每个KV元素分配固定窗口,再通过GroupByKey实现同窗口+同sensorId的数据聚合:

.apply("Apply Fixed Window", Window.into(FixedWindows.of(Duration.standardMinutes(options.getWindowSize()))))
.apply("Group By SensorId", GroupByKey.create())

注意:窗口必须在GroupByKey之前应用,Beam会先将元素分配到对应窗口,再按key分组每个窗口内的元素,确保每个sensorId在独立窗口内聚合。

4. 自定义持久化逻辑

为了让每个传感器的窗口数据独立存储,需要自定义写入逻辑,将sensorId和窗口时间范围融入文件名:

import org.apache.beam.sdk.io.FileIO;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;

public class WriteSensorWindowToGCS extends PTransform<PCollection<KV<String, Iterable<SensorData>>>, PCollection<FileIO.Write.FileResult>> {
    private final String baseOutputPath;

    public WriteSensorWindowToGCS(String baseOutputPath) {
        this.baseOutputPath = baseOutputPath;
    }

    @Override
    public PCollection<FileIO.Write.FileResult> expand(PCollection<KV<String, Iterable<SensorData>>> input) {
        return input
            .apply("Prepare File Outputs", ParDo.of(new DoFn<KV<String, Iterable<SensorData>>, FileIO.Write.FileSpec<String>>() {
                @ProcessElement
                public void processElement(ProcessContext c, BoundedWindow window) {
                    String sensorId = c.element().getKey();
                    // 生成窗口时间戳字符串,用于区分不同窗口
                    String windowStart = String.valueOf(window.start().getMillis());
                    String windowEnd = String.valueOf(window.end().getMillis());
                    // 构建输出路径:basePath/sensorId_窗口起始_窗口结束.txt
                    String filePath = String.format("%s/%s_%s_%s.txt", baseOutputPath, sensorId, windowStart, windowEnd);
                    
                    // 将分组后的传感器数据拼接为字符串
                    StringBuilder content = new StringBuilder();
                    Gson gson = new Gson();
                    for (SensorData data : c.element().getValue()) {
                        content.append(gson.toJson(data)).append("\n");
                    }
                    
                    c.output(FileIO.Write.to(filePath).withContent(content.toString()));
                }
            }))
            .apply(FileIO.write());
    }
}

// 在Pipeline中替换原有写入步骤
.apply("Write Sensor Window Data", new WriteSensorWindowToGCS(options.getOutput()))

完整Pipeline示例

import com.google.gson.Gson;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.FileIO;
import org.apache.beam.sdk.io.PubsubIO;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.GroupByKey;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.TypeDescriptors;
import org.joda.time.Duration;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
import org.apache.beam.sdk.transforms.windowing.FixedWindows;
import org.apache.beam.sdk.transforms.windowing.Window;

public class SensorWindowPipeline {
    public static void main(String[] args) {
        PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create();
        Pipeline pipeline = Pipeline.create(options);

        pipeline
            .apply("Read PubSub Messages", PubsubIO.readStrings().fromTopic(options.getInputTopic()))
            .apply("Parse JSON to KV Pair", MapElements.into(TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptors.of(SensorData.class)))
                .via(message -> {
                    Gson gson = new Gson();
                    SensorData data = gson.fromJson(message, SensorData.class);
                    return KV.of(data.getSensorId(), data);
                }))
            .apply("Apply Fixed Window", Window.into(FixedWindows.of(Duration.standardMinutes(options.getWindowSize()))))
            .apply("Group By SensorId", GroupByKey.create())
            .apply("Write to GCS", new WriteSensorWindowToGCS(options.getOutput()));

        pipeline.run().waitUntilFinish();
    }

    // 传感器数据模型
    public static class SensorData {
        @SerializedName("deviceId")
        private String deviceId;
        @SerializedName("value")
        private double value;
        @SerializedName("sensorId")
        private String sensorId;
        @SerializedName("timeInNanoSeconds")
        private long timeInNanoSeconds;
        @SerializedName("receivedTimeInNanoSeconds")
        private long receivedTimeInNanoSeconds;

        public SensorData() {}

        // Getter和Setter方法
        public String getDeviceId() { return deviceId; }
        public void setDeviceId(String deviceId) { this.deviceId = deviceId; }
        public double getValue() { return value; }
        public void setValue(double value) { this.value = value; }
        public String getSensorId() { return sensorId; }
        public void setSensorId(String sensorId) { this.sensorId = sensorId; }
        public long getTimeInNanoSeconds() { return timeInNanoSeconds; }
        public void setTimeInNanoSeconds(long timeInNanoSeconds) { this.timeInNanoSeconds = timeInNanoSeconds; }
        public long getReceivedTimeInNanoSeconds() { return receivedTimeInNanoSeconds; }
        public void setReceivedTimeInNanoSeconds(long receivedTimeInNanoSeconds) { this.receivedTimeInNanoSeconds = receivedTimeInNanoSeconds; }
    }

    // 自定义GCS写入组件
    public static class WriteSensorWindowToGCS extends PTransform<PCollection<KV<String, Iterable<SensorData>>>, PCollection<FileIO.Write.FileResult>> {
        private final String baseOutputPath;

        public WriteSensorWindowToGCS(String baseOutputPath) {
            this.baseOutputPath = baseOutputPath;
        }

        @Override
        public PCollection<FileIO.Write.FileResult> expand(PCollection<KV<String, Iterable<SensorData>>> input) {
            return input
                .apply("Prepare File Outputs", ParDo.of(new DoFn<KV<String, Iterable<SensorData>>, FileIO.Write.FileSpec<String>>() {
                    @ProcessElement
                    public void processElement(ProcessContext c, BoundedWindow window) {
                        String sensorId = c.element().getKey();
                        String windowStart = String.valueOf(window.start().getMillis());
                        String windowEnd = String.valueOf(window.end().getMillis());
                        String filePath = String.format("%s/%s_%s_%s.txt", baseOutputPath, sensorId, windowStart, windowEnd);
                        
                        StringBuilder content = new StringBuilder();
                        Gson gson = new Gson();
                        for (SensorData data : c.element().getValue()) {
                            content.append(gson.toJson(data)).append("\n");
                        }
                        
                        c.output(FileIO.Write.to(filePath).withContent(content.toString()));
                    }
                }))
                .apply(FileIO.write());
        }
    }
}

关键注意事项

  • 依赖配置:确保项目中添加Gson依赖(Maven:com.google.code.gson:gson:2.10.1)。
  • 窗口与分组顺序:必须先应用窗口再执行GroupByKey,否则分组会跨窗口,不符合需求。
  • 输出隔离:通过文件名包含sensorId和窗口时间,避免不同传感器或窗口的数据混存,方便后续检索。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 15:20:25