muralibasani commented on code in PR #22901:
URL: https://github.com/apache/kafka/pull/22901#discussion_r3669889924
##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/HeadersStoreUpgradeIntegrationTest.java:
##########
@@ -135,6 +137,117 @@ private void buildAndStart(final StoreBuilder<?>
storeBuilder,
IntegrationTestUtils.startApplicationAndWaitUntilRunning(kafkaStreams);
}
+ /**
+ * Produces a single record (no explicit timestamp) into the input stream.
+ */
+ private <K, V> void produce(final K key, final V value) {
+ IntegrationTestUtils.produceKeyValuesSynchronously(
+ inputStream,
+ singletonList(KeyValue.pair(key, value)),
+ TestUtils.producerConfig(CLUSTER.bootstrapServers(),
StringSerializer.class, StringSerializer.class),
+ CLUSTER.time,
+ false);
+ }
+
+ /**
+ * Produces a single record with an explicit timestamp into the input
stream.
+ */
+ private <K, V> void produce(final K key, final V value, final long
timestamp) {
+ IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp(
+ inputStream,
+ singletonList(KeyValue.pair(key, value)),
+ TestUtils.producerConfig(CLUSTER.bootstrapServers(),
StringSerializer.class, StringSerializer.class),
+ timestamp,
+ false);
+ }
+
+ /**
+ * Produces a single record with an explicit timestamp and headers into
the input stream.
+ */
+ private <K, V> void produce(final K key, final V value, final long
timestamp, final Headers headers) {
Review Comment:
Done, switched to produce("key1", "value1", baseTime + 100, headers)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]