基于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.
transform() Core Approach
We’ll split the process into two sequential transform() stages:
- Daily Aggregation: Accumulate values per key and calendar day, emitting daily totals at the end of each day.
- 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 anIterable<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

