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

基于Kafka Streams实现日度与月度数据转换聚合

Alright, let's walk through how to build this two-step aggregation (daily rollup first, then monthly cumulative aggregation) using Kafka Streams 0.10.1's transform() method. Since this version supports branching with array/iterable returns, we can leverage that (along with state stores) to handle intermediate and final aggregates smoothly.

Implementing Daily → Monthly Aggregation with Kafka Streams 0.10.1 transform()

Core Approach

We’ll split the process into two sequential transform() stages:

  1. Daily Aggregation: Accumulate values per key and calendar day, emitting daily totals at the end of each day.
  2. Monthly Rollup: Take the daily totals and accumulate them per key and calendar month—either emitting running totals in real-time or final monthly totals at month-end.

Both stages use persistent key-value state stores to maintain intermediate aggregates, and we’ll use punctuate() to trigger periodic result emission.


Step 1: Configure State Stores

First, define persistent state stores to hold our daily and monthly aggregates. These will be attached to our topology to survive application restarts:

// Build state stores for daily and monthly aggregation
StateStore dailyAggStore = Stores.keyValueStoreBuilder(
    Stores.persistentKeyValueStore("daily-agg-state"),
    Serdes.String(),
    Serdes.Double()
).build();

StateStore monthlyAggStore = Stores.keyValueStoreBuilder(
    Stores.persistentKeyValueStore("monthly-agg-state"),
    Serdes.String(),
    Serdes.Double()
).build();

// Add stores to the StreamsBuilder
StreamsBuilder builder = new StreamsBuilder();
builder.addStateStore(dailyAggStore);
builder.addStateStore(monthlyAggStore);

Step 2: Implement Daily Aggregation Transformer

This transformer accumulates values per key and day, then emits the daily total once the day ends. We use punctuate() to schedule a daily check for expired aggregates:

public class DailyAggTransformer implements Transformer<String, Double, KeyValue<String, Double>> {
    private KeyValueStore<String, Double> stateStore;
    private ProcessorContext context;

    @Override
    public void init(ProcessorContext context) {
        this.context = context;
        // Retrieve the configured state store
        stateStore = (KeyValueStore<String, Double>) context.getStateStore("daily-agg-state");
        // Schedule a daily punctuate (24 hours in milliseconds)
        context.schedule(86400000L);
    }

    @Override
    public KeyValue<String, Double> transform(String key, Double value) {
        // Derive date from record timestamp (use event time if available)
        LocalDateTime eventTime = LocalDateTime.ofInstant(
            Instant.ofEpochMilli(context.timestamp()),
            ZoneId.systemDefault()
        );
        String dateKey = String.format("%s-%s", key, eventTime.format(DateTimeFormatter.ISO_LOCAL_DATE));

        // Update daily aggregate in state
        Double currentTotal = stateStore.get(dateKey);
        Double newTotal = (currentTotal == null) ? value : currentTotal + value;
        stateStore.put(dateKey, newTotal);

        // Return null here—we'll emit results via punctuate()
        return null;
    }

    @Override
    public void punctuate(long timestamp) {
        // Emit aggregates for the previous day
        LocalDateTime yesterday = LocalDateTime.ofInstant(
            Instant.ofEpochMilli(timestamp),
            ZoneId.systemDefault()
        ).minusDays(1);
        String yesterdayStr = yesterday.format(DateTimeFormatter.ISO_LOCAL_DATE);

        // Iterate through state and forward relevant entries
        try (KeyValueIterator<String, Double> iterator = stateStore.all()) {
            while (iterator.hasNext()) {
                KeyValue<String, Double> entry = iterator.next();
                if (entry.key.endsWith("-" + yesterdayStr)) {
                    String originalKey = entry.key.substring(0, entry.key.lastIndexOf("-"));
                    // Forward the daily total to the next stage
                    context.forward(originalKey, entry.value);
                    // Clean up old state to save space
                    stateStore.delete(entry.key);
                }
            }
        }
    }

    @Override
    public void close() {
        // No cleanup needed for this example
    }
}

Step 3: Implement Monthly Rollup Transformer

This transformer takes daily totals and accumulates them into monthly totals. We can choose to emit running monthly totals immediately, or final totals at month-end via punctuate():

public class MonthlyRollupTransformer implements Transformer<String, Double, KeyValue<String, Double>> {
    private KeyValueStore<String, Double> stateStore;
    private ProcessorContext context;

    @Override
    public void init(ProcessorContext context) {
        this.context = context;
        stateStore = (KeyValueStore<String, Double>) context.getStateStore("monthly-agg-state");
        // Optional: Schedule monthly punctuate for final total emission (30 days in ms)
        context.schedule(2592000000L);
    }

    @Override
    public KeyValue<String, Double> transform(String key, Double dailyTotal) {
        // Derive month from record timestamp
        LocalDateTime eventTime = LocalDateTime.ofInstant(
            Instant.ofEpochMilli(context.timestamp()),
            ZoneId.systemDefault()
        );
        String monthKey = String.format("%s-%s", key, eventTime.format(DateTimeFormatter.ofPattern("yyyy-MM")));

        // Update monthly aggregate in state
        Double currentMonthlyTotal = stateStore.get(monthKey);
        Double newMonthlyTotal = (currentMonthlyTotal == null) ? dailyTotal : currentMonthlyTotal + dailyTotal;
        stateStore.put(monthKey, newMonthlyTotal);

        // Return the running monthly total to emit immediately
        return KeyValue.pair(key, newMonthlyTotal);
    }

    @Override
    public void punctuate(long timestamp) {
        // Emit final monthly totals for the previous month and clean up state
        LocalDateTime lastMonth = LocalDateTime.ofInstant(
            Instant.ofEpochMilli(timestamp),
            ZoneId.systemDefault()
        ).minusMonths(1);
        String lastMonthStr = lastMonth.format(DateTimeFormatter.ofPattern("yyyy-MM"));

        try (KeyValueIterator<String, Double> iterator = stateStore.all()) {
            while (iterator.hasNext()) {
                KeyValue<String, Double> entry = iterator.next();
                if (entry.key.endsWith("-" + lastMonthStr)) {
                    String originalKey = entry.key.substring(0, entry.key.lastIndexOf("-"));
                    // Forward final monthly total with a header for identification
                    context.forward(
                        originalKey,
                        entry.value,
                        To.all().withHeader("aggregate-type", "final-monthly".getBytes())
                    );
                    stateStore.delete(entry.key);
                }
            }
        }
    }

    @Override
    public void close() {
        // No cleanup needed for this example
    }
}

Step 4: Wire Up the Topology

Connect the input stream to the two transformers, then output results to target topics:

// Create input stream from source topic
KStream<String, Double> inputStream = builder.stream("input-topic");

// Apply daily aggregation transform
KStream<String, Double> dailyAggStream = inputStream.transform(
    DailyAggTransformer::new,
    "daily-agg-state" // Bind to our state store
);

// Apply monthly rollup transform
KStream<String, Double> monthlyAggStream = dailyAggStream.transform(
    MonthlyRollupTransformer::new,
    "monthly-agg-state"
);

// Output results to target topics
dailyAggStream.to("daily-agg-topic");
monthlyAggStream.to("monthly-agg-topic");

// Build and start the streams application
KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfig);
streams.start();

Key Notes for Kafka Streams 0.10.1

  • State Store Persistence: Using persistent stores ensures aggregates survive application restarts and failures.
  • Punctuate Timing: The schedule() interval uses processing time—if you need event-time based triggering, implement custom logic to track event days/months instead.
  • Iterable Returns: If you need to emit multiple records from a single transform() call, return an Iterable<KeyValue<K, V>> instead of a single entry. For example, emitting both a daily update and a running total in one step.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:38:08