This is an automated email from the ASF dual-hosted git repository.

frankvicky pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new ebac341b28d KAFKA-20278 : streams - Refactor setup pattern, PART-1 
(#21754)
ebac341b28d is described below

commit ebac341b28d4224c296ada31eb45122176e8b27b
Author: Murali Basani <[email protected]>
AuthorDate: Tue Jul 21 15:37:46 2026 +0200

    KAFKA-20278 : streams - Refactor setup pattern, PART-1 (#21754)
    
    Note : This pr is just part 1 of KAFKA-20278
    
    Refactor setup pattern : multiple things - create builder, add state
    store, add stream, add processor, build topology, create KafkaStreams,
    start
    
    Introduced two reusable helper methods
    - buildAndStart
    - assertDowngradeThrowsProcessorStateException
    
    Reviewers: Alieh Saeedi <[email protected]>, TengYao Chi
     <[email protected]>
---
 .../HeadersStoreUpgradeIntegrationTest.java        | 675 ++++++++-------------
 1 file changed, 245 insertions(+), 430 deletions(-)

diff --git 
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/HeadersStoreUpgradeIntegrationTest.java
 
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/HeadersStoreUpgradeIntegrationTest.java
index bd5d5b87a56..f1745709ad7 100644
--- 
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/HeadersStoreUpgradeIntegrationTest.java
+++ 
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/HeadersStoreUpgradeIntegrationTest.java
@@ -33,6 +33,7 @@ import org.apache.kafka.streams.kstream.Windowed;
 import org.apache.kafka.streams.kstream.internals.SessionWindow;
 import org.apache.kafka.streams.processor.api.Processor;
 import org.apache.kafka.streams.processor.api.ProcessorContext;
+import org.apache.kafka.streams.processor.api.ProcessorSupplier;
 import org.apache.kafka.streams.processor.api.Record;
 import org.apache.kafka.streams.state.AggregationWithHeaders;
 import org.apache.kafka.streams.state.KeyValueIterator;
@@ -43,6 +44,7 @@ import org.apache.kafka.streams.state.ReadOnlySessionStore;
 import org.apache.kafka.streams.state.ReadOnlyWindowStore;
 import org.apache.kafka.streams.state.SessionStore;
 import org.apache.kafka.streams.state.SessionStoreWithHeaders;
+import org.apache.kafka.streams.state.StoreBuilder;
 import org.apache.kafka.streams.state.Stores;
 import org.apache.kafka.streams.state.TimestampedKeyValueStore;
 import org.apache.kafka.streams.state.TimestampedKeyValueStoreWithHeaders;
@@ -121,6 +123,56 @@ public class HeadersStoreUpgradeIntegrationTest {
         return streamsConfiguration;
     }
 
+    private void buildAndStart(final StoreBuilder<?> storeBuilder,
+                               final ProcessorSupplier<String, String, Void, 
Void> processorSupplier,
+                               final String storeName,
+                               final Properties props) throws Exception {
+        final StreamsBuilder builder = new StreamsBuilder();
+        builder.addStateStore(storeBuilder)
+            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
+            .process(processorSupplier, storeName);
+        kafkaStreams = new KafkaStreams(builder.build(), props);
+        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+    }
+
+    private void assertDowngradeThrowsProcessorStateException(
+            final String downgradeTarget,
+            final StoreBuilder<?> storeBuilder,
+            final ProcessorSupplier<String, String, Void, Void> 
processorSupplier,
+            final String storeName,
+            final Properties props) {
+        boolean exceptionThrown = false;
+        try {
+            buildAndStart(storeBuilder, processorSupplier, storeName, props);
+        } catch (final Exception e) {
+            Throwable cause = e;
+            while (cause != null) {
+                if (cause instanceof ProcessorStateException &&
+                    cause.getMessage() != null &&
+                    cause.getMessage().contains("headers-aware") &&
+                    cause.getMessage().contains("Downgrade")) {
+                    exceptionThrown = true;
+                    break;
+                }
+                cause = cause.getCause();
+            }
+            if (!exceptionThrown) {
+                throw new AssertionError(
+                    "Expected ProcessorStateException about downgrade " + 
downgradeTarget
+                        + " not being supported, but got: " + e.getMessage(), 
e);
+            }
+        } finally {
+            if (kafkaStreams != null) {
+                kafkaStreams.close(Duration.ofSeconds(30L));
+            }
+        }
+        if (!exceptionThrown) {
+            throw new AssertionError(
+                "Expected ProcessorStateException to be thrown when attempting 
to downgrade "
+                    + downgradeTarget + " from headers-aware store");
+        }
+    }
+
     @AfterEach
     public void shutdown() {
         if (kafkaStreams != null) {
@@ -140,19 +192,14 @@ public class HeadersStoreUpgradeIntegrationTest {
     }
 
     private void 
shouldMigrateTimestampedKeyValueStoreToTimestampedKeyValueStoreWithHeadersUsingPapi(final
 boolean persistentStore) throws Exception {
-        final StreamsBuilder streamsBuilderForOldStore = new StreamsBuilder();
-
-        streamsBuilderForOldStore.addStateStore(
-                Stores.timestampedKeyValueStoreBuilder(
-                    persistentStore ? 
Stores.persistentTimestampedKeyValueStore(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueProcessor::new, STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(streamsBuilderForOldStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        buildAndStart(
+            Stores.timestampedKeyValueStoreBuilder(
+                persistentStore ? 
Stores.persistentTimestampedKeyValueStore(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueProcessor::new, STORE_NAME, props);
 
         processKeyValueAndVerifyTimestampedValue("key1", "value1", 11L);
         processKeyValueAndVerifyTimestampedValue("key2", "value2", 22L);
@@ -161,18 +208,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.close();
         kafkaStreams = null;
 
-        final StreamsBuilder streamsBuilderForNewStore = new StreamsBuilder();
-
-        streamsBuilderForNewStore.addStateStore(
-                Stores.timestampedKeyValueStoreWithHeadersBuilder(
-                    persistentStore ? 
Stores.persistentTimestampedKeyValueStoreWithHeaders(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(streamsBuilderForNewStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedKeyValueStoreWithHeadersBuilder(
+                persistentStore ? 
Stores.persistentTimestampedKeyValueStoreWithHeaders(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME, props);
 
         // Verify legacy data can be read with empty headers
         verifyLegacyValuesWithEmptyHeaders("key1", "value1", 11L);
@@ -191,19 +232,14 @@ public class HeadersStoreUpgradeIntegrationTest {
 
     @Test
     public void 
shouldProxyTimestampedKeyValueStoreToTimestampedKeyValueStoreWithHeadersUsingPapi()
 throws Exception {
-        final StreamsBuilder streamsBuilderForOldStore = new StreamsBuilder();
-
-        streamsBuilderForOldStore.addStateStore(
-                Stores.timestampedKeyValueStoreBuilder(
-                    Stores.persistentTimestampedKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueProcessor::new, STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(streamsBuilderForOldStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        buildAndStart(
+            Stores.timestampedKeyValueStoreBuilder(
+                Stores.persistentTimestampedKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueProcessor::new, STORE_NAME, props);
 
         processKeyValueAndVerifyTimestampedValue("key1", "value1", 11L);
         processKeyValueAndVerifyTimestampedValue("key2", "value2", 22L);
@@ -212,20 +248,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.close();
         kafkaStreams = null;
 
-
-
-        final StreamsBuilder streamsBuilderForNewStore = new StreamsBuilder();
-
-        streamsBuilderForNewStore.addStateStore(
-                Stores.timestampedKeyValueStoreWithHeadersBuilder(
-                    Stores.persistentTimestampedKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(streamsBuilderForNewStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedKeyValueStoreWithHeadersBuilder(
+                Stores.persistentTimestampedKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME, props);
 
         // Verify legacy data can be read with empty headers
         verifyLegacyValuesWithEmptyHeaders("key1", "value1", 11L);
@@ -254,19 +282,14 @@ public class HeadersStoreUpgradeIntegrationTest {
     }
 
     private void 
shouldMigratePlainKeyValueStoreToTimestampedKeyValueStoreWithHeadersUsingPapi(final
 boolean persistentStore) throws Exception {
-        final StreamsBuilder streamsBuilderForOldStore = new StreamsBuilder();
-
-        streamsBuilderForOldStore.addStateStore(
-                Stores.keyValueStoreBuilder(
-                    persistentStore ? 
Stores.persistentKeyValueStore(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(KeyValueProcessor::new, STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(streamsBuilderForOldStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        buildAndStart(
+            Stores.keyValueStoreBuilder(
+                persistentStore ? Stores.persistentKeyValueStore(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            KeyValueProcessor::new, STORE_NAME, props);
 
         processKeyValueAndVerifyValue("key1", "value1");
         final long lastUpdateKeyOne = persistentStore ? -1L : 
CLUSTER.time.milliseconds() - 1L;
@@ -280,18 +303,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.close();
         kafkaStreams = null;
 
-        final StreamsBuilder streamsBuilderForNewStore = new StreamsBuilder();
-
-        streamsBuilderForNewStore.addStateStore(
-                Stores.timestampedKeyValueStoreWithHeadersBuilder(
-                    persistentStore ? 
Stores.persistentTimestampedKeyValueStoreWithHeaders(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(streamsBuilderForNewStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedKeyValueStoreWithHeadersBuilder(
+                persistentStore ? 
Stores.persistentTimestampedKeyValueStoreWithHeaders(STORE_NAME) : 
Stores.inMemoryKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME, props);
 
         // Verify legacy data can be read with empty headers and timestamp
         verifyLegacyValuesWithEmptyHeaders("key1", "value1", lastUpdateKeyOne);
@@ -310,19 +327,14 @@ public class HeadersStoreUpgradeIntegrationTest {
 
     @Test
     public void 
shouldProxyPlainKeyValueStoreToTimestampedKeyValueStoreWithHeadersUsingPapi() 
throws Exception {
-        final StreamsBuilder streamsBuilderForOldStore = new StreamsBuilder();
-
-        streamsBuilderForOldStore.addStateStore(
-                Stores.keyValueStoreBuilder(
-                    Stores.persistentKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(KeyValueProcessor::new, STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(streamsBuilderForOldStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        buildAndStart(
+            Stores.keyValueStoreBuilder(
+                Stores.persistentKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            KeyValueProcessor::new, STORE_NAME, props);
 
         processKeyValueAndVerifyValue("key1", "value1");
         processKeyValueAndVerifyValue("key2", "value2");
@@ -331,20 +343,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.close();
         kafkaStreams = null;
 
-
-
-        final StreamsBuilder streamsBuilderForNewStore = new StreamsBuilder();
-
-        streamsBuilderForNewStore.addStateStore(
-                Stores.timestampedKeyValueStoreWithHeadersBuilder(
-                    Stores.persistentKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(streamsBuilderForNewStore.build(), 
props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedKeyValueStoreWithHeadersBuilder(
+                Stores.persistentKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME, props);
 
         // Verify legacy data can be read with empty headers
         verifyLegacyValuesWithEmptyHeaders("key1", "value1", -1L);
@@ -617,21 +621,17 @@ public class HeadersStoreUpgradeIntegrationTest {
     }
 
     private void 
shouldMigratePlainWindowStoreToTimestampedWindowStoreWithHeaders(final boolean 
persistentStore) throws Exception {
-        // Run with old plain WindowStore
-        final StreamsBuilder oldBuilder = new StreamsBuilder();
-        oldBuilder.addStateStore(
-                Stores.windowStoreBuilder(
-                    persistentStore
-                        ? Stores.persistentWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
-                        : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(PlainWindowedProcessor::new, WINDOW_STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(oldBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        // Run with old plain WindowStore
+        buildAndStart(
+            Stores.windowStoreBuilder(
+                persistentStore
+                    ? Stores.persistentWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
+                    : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            PlainWindowedProcessor::new, WINDOW_STORE_NAME, props);
 
         final long baseTime = CLUSTER.time.milliseconds();
         processPlainWindowedKeyValueAndVerify("key1", "value1", baseTime + 
100);
@@ -642,19 +642,14 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams = null;
 
         // Restart with TimestampedWindowStoreWithHeaders
-        final StreamsBuilder newBuilder = new StreamsBuilder();
-        newBuilder.addStateStore(
-                Stores.timestampedWindowStoreWithHeadersBuilder(
-                    persistentStore
-                        ? 
Stores.persistentTimestampedWindowStoreWithHeaders(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
-                        : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedWithHeadersProcessor::new, 
WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(newBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedWindowStoreWithHeadersBuilder(
+                persistentStore
+                    ? 
Stores.persistentTimestampedWindowStoreWithHeaders(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
+                    : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedWithHeadersProcessor::new, WINDOW_STORE_NAME, 
props);
 
         verifyPlainWindowValueWithEmptyHeadersAndTimestamp("key1", "value1", 
baseTime + 100, persistentStore ? -1L : baseTime + 100);
         verifyPlainWindowValueWithEmptyHeadersAndTimestamp("key2", "value2", 
baseTime + 200, persistentStore ? -1L : baseTime + 200);
@@ -672,18 +667,14 @@ public class HeadersStoreUpgradeIntegrationTest {
 
     @Test
     public void 
shouldProxyPlainWindowStoreToTimestampedWindowStoreWithHeaders() throws 
Exception {
-        final StreamsBuilder oldBuilder = new StreamsBuilder();
-        oldBuilder.addStateStore(
-                Stores.windowStoreBuilder(
-                    Stores.persistentWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(PlainWindowedProcessor::new, WINDOW_STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(oldBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        buildAndStart(
+            Stores.windowStoreBuilder(
+                Stores.persistentWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            PlainWindowedProcessor::new, WINDOW_STORE_NAME, props);
 
         final long baseTime = CLUSTER.time.milliseconds();
         processPlainWindowedKeyValueAndVerify("key1", "value1", baseTime + 
100);
@@ -694,17 +685,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams = null;
 
         // Restart with headers-aware builder but non-headers supplier 
(proxy/adapter mode)
-        final StreamsBuilder newBuilder = new StreamsBuilder();
-        newBuilder.addStateStore(
-                Stores.timestampedWindowStoreWithHeadersBuilder(
-                    Stores.persistentWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),  // 
non-headers supplier!
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedWithHeadersProcessor::new, 
WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(newBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedWindowStoreWithHeadersBuilder(
+                Stores.persistentWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedWithHeadersProcessor::new, WINDOW_STORE_NAME, 
props);
 
         verifyPlainWindowValueWithEmptyHeadersAndTimestamp("key1", "value1", 
baseTime + 100, -1L);
         verifyPlainWindowValueWithEmptyHeadersAndTimestamp("key2", "value2", 
baseTime + 200, -1L);
@@ -738,21 +724,17 @@ public class HeadersStoreUpgradeIntegrationTest {
      * This is a true migration where both supplier and builder are upgraded.
      */
     private void 
shouldMigrateTimestampedWindowStoreToTimestampedWindowStoreWithHeaders(final 
boolean persistentStore) throws Exception {
-        // Phase 1: Run with old TimestampedWindowStore
-        final StreamsBuilder oldBuilder = new StreamsBuilder();
-        oldBuilder.addStateStore(
-                Stores.timestampedWindowStoreBuilder(
-                    persistentStore
-                        ? 
Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
-                        : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedProcessor::new, WINDOW_STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(oldBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        // Phase 1: Run with old TimestampedWindowStore
+        buildAndStart(
+            Stores.timestampedWindowStoreBuilder(
+                persistentStore
+                    ? 
Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
+                    : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedProcessor::new, WINDOW_STORE_NAME, props);
 
         final long baseTime = CLUSTER.time.milliseconds();
         processWindowedKeyValueAndVerifyTimestamped("key1", "value1", baseTime 
+ 100);
@@ -762,19 +744,14 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.close();
         kafkaStreams = null;
 
-        final StreamsBuilder newBuilder = new StreamsBuilder();
-        newBuilder.addStateStore(
-                Stores.timestampedWindowStoreWithHeadersBuilder(
-                    persistentStore
-                        ? 
Stores.persistentTimestampedWindowStoreWithHeaders(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
-                        : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedWithHeadersProcessor::new, 
WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(newBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedWindowStoreWithHeadersBuilder(
+                persistentStore
+                    ? 
Stores.persistentTimestampedWindowStoreWithHeaders(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false)
+                    : Stores.inMemoryWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedWithHeadersProcessor::new, WINDOW_STORE_NAME, 
props);
 
         verifyWindowValueWithEmptyHeaders("key1", "value1", baseTime + 100);
         verifyWindowValueWithEmptyHeaders("key2", "value2", baseTime + 200);
@@ -792,18 +769,14 @@ public class HeadersStoreUpgradeIntegrationTest {
 
     @Test
     public void 
shouldProxyTimestampedWindowStoreToTimestampedWindowStoreWithHeaders() throws 
Exception {
-        final StreamsBuilder oldBuilder = new StreamsBuilder();
-        oldBuilder.addStateStore(
-                Stores.timestampedWindowStoreBuilder(
-                    Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedProcessor::new, WINDOW_STORE_NAME);
-
         final Properties props = props();
-        kafkaStreams = new KafkaStreams(oldBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+
+        buildAndStart(
+            Stores.timestampedWindowStoreBuilder(
+                Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedProcessor::new, WINDOW_STORE_NAME, props);
 
         final long baseTime = CLUSTER.time.milliseconds();
         processWindowedKeyValueAndVerifyTimestamped("key1", "value1", baseTime 
+ 100);
@@ -814,17 +787,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams = null;
 
         // Restart with headers-aware builder but non-headers supplier 
(proxy/adapter mode)
-        final StreamsBuilder newBuilder = new StreamsBuilder();
-        newBuilder.addStateStore(
-                Stores.timestampedWindowStoreWithHeadersBuilder(
-                    Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),  // 
non-headers supplier!
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedWithHeadersProcessor::new, 
WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(newBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedWindowStoreWithHeadersBuilder(
+                Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME, 
Duration.ofMillis(RETENTION_MS), Duration.ofMillis(WINDOW_SIZE_MS), false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedWithHeadersProcessor::new, WINDOW_STORE_NAME, 
props);
 
         verifyWindowValueWithEmptyHeaders("key1", "value1", baseTime + 100);
         verifyWindowValueWithEmptyHeaders("key2", "value2", baseTime + 200);
@@ -1165,44 +1133,13 @@ public class HeadersStoreUpgradeIntegrationTest {
         setupAndPopulateKeyValueStoreWithHeaders(props);
         kafkaStreams = null;
 
-        // Attempt to downgrade to plain key-value store
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.keyValueStoreBuilder(
-                    Stores.persistentKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(KeyValueProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-
-        boolean exceptionThrown = false;
-        try {
-            
IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
-        } catch (final Exception e) {
-            Throwable cause = e;
-            while (cause != null) {
-                if (cause instanceof ProcessorStateException &&
-                    cause.getMessage() != null &&
-                    cause.getMessage().contains("headers-aware") &&
-                    cause.getMessage().contains("Downgrade")) {
-                    exceptionThrown = true;
-                    break;
-                }
-                cause = cause.getCause();
-            }
-
-            if (!exceptionThrown) {
-                throw new AssertionError("Expected ProcessorStateException 
about downgrade not being supported, but got: " + e.getMessage(), e);
-            }
-        } finally {
-            kafkaStreams.close(Duration.ofSeconds(30L));
-        }
-
-        if (!exceptionThrown) {
-            throw new AssertionError("Expected ProcessorStateException to be 
thrown when attempting to downgrade from headers-aware to plain key-value 
store");
-        }
+        assertDowngradeThrowsProcessorStateException(
+            "to plain key-value store",
+            Stores.keyValueStoreBuilder(
+                Stores.persistentKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            KeyValueProcessor::new, STORE_NAME, props);
     }
 
     @Test
@@ -1213,17 +1150,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.cleanUp(); // Delete local state
         kafkaStreams = null;
 
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.keyValueStoreBuilder(
-                    Stores.persistentKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(KeyValueProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.keyValueStoreBuilder(
+                Stores.persistentKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            KeyValueProcessor::new, STORE_NAME, props);
 
         processKeyValueAndVerifyValue("key3", "value3");
         processKeyValueAndVerifyValue("key4", "value4");
@@ -1237,44 +1169,13 @@ public class HeadersStoreUpgradeIntegrationTest {
         setupAndPopulateKeyValueStoreWithHeaders(props);
         kafkaStreams = null;
 
-        // Attempt to downgrade to non-headers key-value store
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.timestampedKeyValueStoreBuilder(
-                    Stores.persistentTimestampedKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-
-        boolean exceptionThrown = false;
-        try {
-            
IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
-        } catch (final Exception e) {
-            Throwable cause = e;
-            while (cause != null) {
-                if (cause instanceof ProcessorStateException &&
-                    cause.getMessage() != null &&
-                    cause.getMessage().contains("headers-aware") &&
-                    cause.getMessage().contains("Downgrade")) {
-                    exceptionThrown = true;
-                    break;
-                }
-                cause = cause.getCause();
-            }
-
-            if (!exceptionThrown) {
-                throw new AssertionError("Expected ProcessorStateException 
about downgrade not being supported, but got: " + e.getMessage(), e);
-            }
-        } finally {
-            kafkaStreams.close(Duration.ofSeconds(30L));
-        }
-
-        if (!exceptionThrown) {
-            throw new AssertionError("Expected ProcessorStateException to be 
thrown when attempting to downgrade from headers-aware to non-headers key-value 
store");
-        }
+        assertDowngradeThrowsProcessorStateException(
+            "to timestamped key-value store",
+            Stores.timestampedKeyValueStoreBuilder(
+                Stores.persistentTimestampedKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueProcessor::new, STORE_NAME, props);
     }
 
     @Test
@@ -1285,17 +1186,12 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.cleanUp(); // Delete local state
         kafkaStreams = null;
 
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.timestampedKeyValueStoreBuilder(
-                    Stores.persistentTimestampedKeyValueStore(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedKeyValueStoreBuilder(
+                Stores.persistentTimestampedKeyValueStore(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueProcessor::new, STORE_NAME, props);
 
         // verify legacy key, values
         verifyLegacyTimestampedValue("key1", "value1", 11L);
@@ -1314,46 +1210,16 @@ public class HeadersStoreUpgradeIntegrationTest {
         setupAndPopulateWindowStoreWithHeaders(props, 
List.of(KeyValue.pair("key1", 100L)));
         kafkaStreams = null;
 
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.windowStoreBuilder(
-                    Stores.persistentWindowStore(WINDOW_STORE_NAME,
-                        Duration.ofMillis(RETENTION_MS),
-                        Duration.ofMillis(WINDOW_SIZE_MS),
-                        false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(PlainWindowedProcessor::new, WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-
-        boolean exceptionThrown = false;
-        try {
-            
IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
-        } catch (final Exception e) {
-            Throwable cause = e;
-            while (cause != null) {
-                if (cause instanceof ProcessorStateException &&
-                    cause.getMessage() != null &&
-                    cause.getMessage().contains("headers-aware") &&
-                    cause.getMessage().contains("Downgrade")) {
-                    exceptionThrown = true;
-                    break;
-                }
-                cause = cause.getCause();
-            }
-
-            if (!exceptionThrown) {
-                throw new AssertionError("Expected ProcessorStateException 
about downgrade not being supported, but got: " + e.getMessage(), e);
-            }
-        } finally {
-            kafkaStreams.close(Duration.ofSeconds(30L));
-        }
-
-        if (!exceptionThrown) {
-            throw new AssertionError("Expected ProcessorStateException to be 
thrown when attempting to downgrade from headers-aware to plain window store");
-        }
+        assertDowngradeThrowsProcessorStateException(
+            "to plain window store",
+            Stores.windowStoreBuilder(
+                Stores.persistentWindowStore(WINDOW_STORE_NAME,
+                    Duration.ofMillis(RETENTION_MS),
+                    Duration.ofMillis(WINDOW_SIZE_MS),
+                    false),
+                Serdes.String(),
+                Serdes.String()),
+            PlainWindowedProcessor::new, WINDOW_STORE_NAME, props);
     }
 
     @Test
@@ -1362,47 +1228,16 @@ public class HeadersStoreUpgradeIntegrationTest {
         setupAndPopulateWindowStoreWithHeaders(props, 
singletonList(KeyValue.pair("key1", 100L)));
         kafkaStreams = null;
 
-        // Attempt to downgrade to non-headers window store
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.timestampedWindowStoreBuilder(
-                    Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME,
-                        Duration.ofMillis(RETENTION_MS),
-                        Duration.ofMillis(WINDOW_SIZE_MS),
-                        false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedProcessor::new, WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-
-        boolean exceptionThrown = false;
-        try {
-            
IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
-        } catch (final Exception e) {
-            Throwable cause = e;
-            while (cause != null) {
-                if (cause instanceof ProcessorStateException &&
-                    cause.getMessage() != null &&
-                    cause.getMessage().contains("headers-aware") &&
-                    cause.getMessage().contains("Downgrade")) {
-                    exceptionThrown = true;
-                    break;
-                }
-                cause = cause.getCause();
-            }
-
-            if (!exceptionThrown) {
-                throw new AssertionError("Expected ProcessorStateException 
about downgrade not being supported, but got: " + e.getMessage(), e);
-            }
-        } finally {
-            kafkaStreams.close(Duration.ofSeconds(30L));
-        }
-
-        if (!exceptionThrown) {
-            throw new AssertionError("Expected ProcessorStateException to be 
thrown when attempting to downgrade from headers-aware to non-headers window 
store");
-        }
+        assertDowngradeThrowsProcessorStateException(
+            "to timestamped window store",
+            Stores.timestampedWindowStoreBuilder(
+                Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME,
+                    Duration.ofMillis(RETENTION_MS),
+                    Duration.ofMillis(WINDOW_SIZE_MS),
+                    false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedProcessor::new, WINDOW_STORE_NAME, props);
     }
 
     @Test
@@ -1413,20 +1248,15 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.cleanUp();
         kafkaStreams = null;
 
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.windowStoreBuilder(
-                    Stores.persistentWindowStore(WINDOW_STORE_NAME,
-                        Duration.ofMillis(RETENTION_MS),
-                        Duration.ofMillis(WINDOW_SIZE_MS),
-                        false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(PlainWindowedProcessor::new, WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.windowStoreBuilder(
+                Stores.persistentWindowStore(WINDOW_STORE_NAME,
+                    Duration.ofMillis(RETENTION_MS),
+                    Duration.ofMillis(WINDOW_SIZE_MS),
+                    false),
+                Serdes.String(),
+                Serdes.String()),
+            PlainWindowedProcessor::new, WINDOW_STORE_NAME, props);
 
         final long newTime = CLUSTER.time.milliseconds();
         processPlainWindowedKeyValueAndVerify("key3", "value3", newTime + 300);
@@ -1443,20 +1273,15 @@ public class HeadersStoreUpgradeIntegrationTest {
         kafkaStreams.cleanUp(); // Delete local state
         kafkaStreams = null;
 
-        final StreamsBuilder downgradedBuilder = new StreamsBuilder();
-        downgradedBuilder.addStateStore(
-                Stores.timestampedWindowStoreBuilder(
-                    Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME,
-                        Duration.ofMillis(RETENTION_MS),
-                        Duration.ofMillis(WINDOW_SIZE_MS),
-                        false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedProcessor::new, WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(downgradedBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedWindowStoreBuilder(
+                Stores.persistentTimestampedWindowStore(WINDOW_STORE_NAME,
+                    Duration.ofMillis(RETENTION_MS),
+                    Duration.ofMillis(WINDOW_SIZE_MS),
+                    false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedProcessor::new, WINDOW_STORE_NAME, props);
 
         final long newTime = CLUSTER.time.milliseconds();
         processWindowedKeyValueAndVerifyTimestamped("key3", "value3", newTime 
+ 300);
@@ -1522,20 +1347,15 @@ public class HeadersStoreUpgradeIntegrationTest {
     }
 
     private long setupWindowStoreWithHeaders(final Properties props) throws 
Exception {
-        final StreamsBuilder headersBuilder = new StreamsBuilder();
-        headersBuilder.addStateStore(
-                Stores.timestampedWindowStoreWithHeadersBuilder(
-                    
Stores.persistentTimestampedWindowStoreWithHeaders(WINDOW_STORE_NAME,
-                        Duration.ofMillis(RETENTION_MS),
-                        Duration.ofMillis(WINDOW_SIZE_MS),
-                        false),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedWindowedWithHeadersProcessor::new, 
WINDOW_STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(headersBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedWindowStoreWithHeadersBuilder(
+                
Stores.persistentTimestampedWindowStoreWithHeaders(WINDOW_STORE_NAME,
+                    Duration.ofMillis(RETENTION_MS),
+                    Duration.ofMillis(WINDOW_SIZE_MS),
+                    false),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedWindowedWithHeadersProcessor::new, WINDOW_STORE_NAME, 
props);
 
         return CLUSTER.time.milliseconds();
     }
@@ -1554,17 +1374,12 @@ public class HeadersStoreUpgradeIntegrationTest {
     }
 
     private void setupAndPopulateKeyValueStoreWithHeaders(final Properties 
props) throws Exception {
-        final StreamsBuilder headersBuilder = new StreamsBuilder();
-        headersBuilder.addStateStore(
-                Stores.timestampedKeyValueStoreWithHeadersBuilder(
-                    
Stores.persistentTimestampedKeyValueStoreWithHeaders(STORE_NAME),
-                    Serdes.String(),
-                    Serdes.String()))
-            .stream(inputStream, Consumed.with(Serdes.String(), 
Serdes.String()))
-            .process(TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME);
-
-        kafkaStreams = new KafkaStreams(headersBuilder.build(), props);
-        IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
+        buildAndStart(
+            Stores.timestampedKeyValueStoreWithHeadersBuilder(
+                
Stores.persistentTimestampedKeyValueStoreWithHeaders(STORE_NAME),
+                Serdes.String(),
+                Serdes.String()),
+            TimestampedKeyValueWithHeadersProcessor::new, STORE_NAME, props);
 
         final Headers headers = new RecordHeaders();
         headers.add("source", "test".getBytes());


Reply via email to