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

Kafka Stream:基于设备时序数据的记录与聚合技术问询

Great question! Let's walk through exactly how to implement aggregation for your device time-series data with Kafka Streams, using your sample records as a guide. I'll cover core setup, code examples, and key best practices to make this robust in production.

1. Clarify Aggregation Goals First

Looking at your data, we can target two common aggregation patterns:

  • All-time aggregation: Calculate cumulative stats (sum, average) for each unique group of device_nm, type, and model
  • Hourly time-windowed aggregation: Since your starttime values are hourly, compute stats per hour for each group (matches your sample's time intervals)
2. Step-by-Step Implementation

2.1 Setup Dependencies & Configuration

First, ensure you have the Kafka Streams and JSON Serde dependencies (for Maven):

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>3.6.1</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams-json</artifactId>
    <version>3.6.1</version>
</dependency>

Configure your Kafka Streams application with core properties:

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "device-aggregation-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, JsonSerde.class);

2.2 Define Data Classes

Create POJOs to map your input JSON records and aggregation results (required for JsonSerde):

// Input device data POJO
public class DeviceData {
    private String device_nm;
    private String type;
    private int mtrc1;
    private int mtrc2;
    private String starttime;
    private String model;

    // Required: no-arg constructor, getters, and setters for JsonSerde
    public DeviceData() {}

    // Getters and setters here
}

// Aggregation result POJO
public class AggregationResult {
    private String groupKey; // Format: device_nm-type-model
    private long totalMtrc1;
    private double avgMtrc2;
    private int recordCount;

    // Required: no-arg constructor, getters, and setters
    public AggregationResult() {}

    // Getters and setters here
}

2.3 Build the Streams Topology

Option 1: All-Time Aggregation

This maintains a persistent state to accumulate stats for each group indefinitely:

StreamsBuilder builder = new StreamsBuilder();

// Read input topic into a stream
KStream<String, DeviceData> deviceStream = builder.stream("device-input-topic");

// Transform to use a composite group key (device_nm-type-model)
KStream<String, DeviceData> keyedStream = deviceStream.map((ignoredKey, data) -> {
    String groupKey = String.join("|", data.getDevice_nm(), data.getType(), data.getModel());
    return KeyValue.pair(groupKey, data);
});

// Aggregate to calculate sum of mtrc1 and average of mtrc2
KTable<String, AggregationResult> allTimeAggregate = keyedStream
        .groupByKey(Grouped.with(Serdes.String(), new JsonSerde<>(DeviceData.class)))
        .aggregate(
                // Initialize empty result
                AggregationResult::new,
                // Update result with each new record
                (groupKey, data, result) -> {
                    result.setGroupKey(groupKey);
                    result.setTotalMtrc1(result.getTotalMtrc1() + data.getMtrc1());
                    result.setRecordCount(result.getRecordCount() + 1);
                    // Calculate rolling average for mtrc2
                    result.setAvgMtrc2(
                        (result.getAvgMtrc2() * (result.getRecordCount() - 1) + data.getMtrc2()) 
                        / result.getRecordCount()
                    );
                    return result;
                },
                // Materialize state to a Kafka-backed store
                Materialized.<String, AggregationResult, KeyValueStore<Bytes, byte[]>>as("all-time-device-stats")
                        .withKeySerde(Serdes.String())
                        .withValueSerde(new JsonSerde<>(AggregationResult.class))
        );

// Write results to an output topic
allTimeAggregate.toStream().to("device-all-time-stats-topic", Produced.with(Serdes.String(), new JsonSerde<>(AggregationResult.class)));

Option 2: Hourly Windowed Aggregation

This groups data into 1-hour windows (matching your starttime intervals) and computes stats per window:

// Extend the topology with windowed aggregation
KTable<Windowed<String>, AggregationResult> hourlyAggregate = keyedStream
        .groupByKey(Grouped.with(Serdes.String(), new JsonSerde<>(DeviceData.class)))
        .windowedBy(TimeWindows.of(Duration.ofHours(1))
                .grace(Duration.ofMinutes(5))) // Allow 5 mins for late-arriving records
        .aggregate(
                AggregationResult::new,
                (groupKey, data, result) -> {
                    result.setGroupKey(groupKey);
                    result.setTotalMtrc1(result.getTotalMtrc1() + data.getMtrc1());
                    result.setRecordCount(result.getRecordCount() + 1);
                    result.setAvgMtrc2(
                        (result.getAvgMtrc2() * (result.getRecordCount() - 1) + data.getMtrc2()) 
                        / result.getRecordCount()
                    );
                    return result;
                },
                Materialized.<String, AggregationResult, WindowStore<Bytes, byte[]>>as("hourly-device-stats")
                        .withKeySerde(Serdes.String())
                        .withValueSerde(new JsonSerde<>(AggregationResult.class))
        );

// Write windowed results (include window start time in output key)
hourlyAggregate.toStream((windowedKey, result) -> 
        windowedKey.key() + "_" + windowedKey.window().start()
).to("device-hourly-stats-topic", Produced.with(Serdes.String(), new JsonSerde<>(AggregationResult.class)));

2.4 Start the Stream & Handle Shutdown

KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

// Add shutdown hook to cleanly close the stream
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
3. Key Technical Guidance & Best Practices
  • State Management: Aggregations are stateful—Kafka Streams uses RocksDB by default for persistent state. Configure state.dir in properties to specify where state is stored (critical for production).
  • Late Data Handling: Use the grace() method on time windows to accept late records (adjust the duration based on your data's expected latency).
  • Key Design: Use a unique composite key (like device_nm|type|model) to avoid grouping conflicts. Avoid using characters that exist in your field values as separators.
  • Monitoring: Enable JMX metrics to track state store size, processing latency, and error rates. Tools like Prometheus/Grafana can help visualize these metrics.
  • Testing: Use TopologyTestDriver to simulate input records and validate aggregation outputs without connecting to a real Kafka cluster.
4. Example Outputs

For your sample data:

  • All-time aggregation for x1|external|t20:
    {
      "groupKey": "x1|external|t20",
      "totalMtrc1": 105,
      "avgMtrc2": 18.0,
      "recordCount": 2
    }
    
  • Hourly aggregation for x1|internal|t40 (2018-05-04 01:00:00 window):
    {
      "groupKey": "x1|internal|t40",
      "totalMtrc1": 35,
      "avgMtrc2": 15.0,
      "recordCount": 1
    }
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:57:52