muralibasani commented on code in PR #22901:
URL: https://github.com/apache/kafka/pull/22901#discussion_r3722859420
##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/HeadersStoreUpgradeIntegrationTest.java:
##########
@@ -814,261 +797,89 @@ public void
shouldProxyTimestampedWindowStoreToTimestampedWindowStoreWithHeaders
private void processPlainWindowedKeyValueAndVerify(final String key,
final String value,
final long timestamp)
throws Exception {
- IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp(
- inputStream,
- List.of(KeyValue.pair(key, value)),
- TestUtils.producerConfig(CLUSTER.bootstrapServers(),
- StringSerializer.class,
- StringSerializer.class),
- timestamp,
- false);
-
- TestUtils.waitForCondition(() -> {
- try {
- final ReadOnlyWindowStore<String, String> store =
- IntegrationTestUtils.getStore(WINDOW_STORE_NAME,
kafkaStreams, QueryableStoreTypes.windowStore());
-
- if (store == null) {
- return false;
- }
-
- final long windowStart = timestamp - (timestamp %
WINDOW_SIZE_MS);
- final String result = store.fetch(key, windowStart);
+ produce(key, value, timestamp);
+ awaitStore(WINDOW_STORE_NAME, QueryableStoreTypes.<String,
String>windowStore(),
+ store -> {
+ final String result = store.fetch(key, windowStart(timestamp));
return result != null && result.equals(value);
- } catch (final Exception e) {
- return false;
- }
- }, 60_000L, "Could not verify plain window value in time.");
+ },
+ "Could not verify plain window value in time.");
}
private void verifyPlainWindowValueWithEmptyHeadersAndTimestamp(final
String key,
final
String value,
final long
windowTimestamp,
final long
expectedTimestamp) throws Exception {
- TestUtils.waitForCondition(() -> {
- try {
- final ReadOnlyWindowStore<String,
ValueTimestampHeaders<String>> store =
- IntegrationTestUtils.getStore(WINDOW_STORE_NAME,
kafkaStreams, QueryableStoreTypes.timestampedWindowStoreWithHeaders());
-
- if (store == null) {
- return false;
- }
-
- final long windowStart = windowTimestamp - (windowTimestamp %
WINDOW_SIZE_MS);
-
- final List<KeyValue<Windowed<String>,
ValueTimestampHeaders<String>>> results = new LinkedList<>();
- try (final KeyValueIterator<Windowed<String>,
ValueTimestampHeaders<String>> iterator = store.all()) {
- while (iterator.hasNext()) {
- final KeyValue<Windowed<String>,
ValueTimestampHeaders<String>> kv = iterator.next();
- if (kv.key.key().equals(key) &&
kv.key.window().start() == windowStart) {
- results.add(kv);
- }
- }
- }
-
- if (results.isEmpty()) {
+ awaitStore(WINDOW_STORE_NAME, QueryableStoreTypes.<String,
String>timestampedWindowStoreWithHeaders(),
+ store -> {
+ final Optional<KeyValue<Windowed<String>,
ValueTimestampHeaders<String>>> result =
+ findWindowedValue(store, key, windowTimestamp);
+ if (result.isEmpty()) {
return false;
}
- final ValueTimestampHeaders<String> result =
results.get(0).value;
- assertNotNull(result, "Result should not be null");
- assertEquals(value, result.value(), "Value should match");
- assertEquals(expectedTimestamp, result.timestamp(), "Timestamp
should be " + expectedTimestamp + " for plain store migration");
+ final ValueTimestampHeaders<String> actual =
result.get().value;
+ assertNotNull(actual, "Stored value should not be null");
+ assertEquals(value, actual.value(), "Value should match");
+ assertEquals(expectedTimestamp, actual.timestamp(), "Timestamp
should be " + expectedTimestamp + " for plain store migration");
// Verify headers exist but are empty (migrated from plain
store without headers or timestamps)
- assertNotNull(result.headers(), "Headers should not be null
for migrated data");
- assertEquals(0, result.headers().toArray().length, "Headers
should be empty for migrated data");
+ assertNotNull(actual.headers(), "Headers should not be null
for migrated data");
+ assertEquals(0, actual.headers().toArray().length, "Headers
should be empty for migrated data");
return true;
- } catch (final Exception e) {
- LOG.error("Error while verifying plain window value with empty
headers and timestamp", e);
- return false;
- }
- }, 60_000L, "Could not verify plain window value with empty headers
and timestamp in time.");
- }
-
- private void processPlainWindowedKeyValueWithHeadersAndVerify(final String
key,
- final String
value,
- final long
timestamp,
- final
Headers headers,
- final
Headers expectedHeaders) throws Exception {
- IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp(
- inputStream,
- List.of(KeyValue.pair(key, value)),
- TestUtils.producerConfig(CLUSTER.bootstrapServers(),
- StringSerializer.class,
- StringSerializer.class),
- headers,
- timestamp,
- false);
-
- TestUtils.waitForCondition(() -> {
- try {
- final ReadOnlyWindowStore<String,
ValueTimestampHeaders<String>> store =
- IntegrationTestUtils.getStore(WINDOW_STORE_NAME,
kafkaStreams, QueryableStoreTypes.timestampedWindowStoreWithHeaders());
-
- if (store == null) {
- return false;
- }
-
- final long windowStart = timestamp - (timestamp %
WINDOW_SIZE_MS);
-
- final List<KeyValue<Windowed<String>,
ValueTimestampHeaders<String>>> results = new LinkedList<>();
- try (final KeyValueIterator<Windowed<String>,
ValueTimestampHeaders<String>> iterator = store.all()) {
- while (iterator.hasNext()) {
- final KeyValue<Windowed<String>,
ValueTimestampHeaders<String>> kv = iterator.next();
- if (kv.key.key().equals(key) &&
kv.key.window().start() == windowStart) {
- results.add(kv);
- }
- }
- }
-
- if (results.isEmpty()) {
- return false;
- }
-
- final ValueTimestampHeaders<String> result =
results.get(0).value;
- // For plain window stores, timestamp is always -1 since it's
not preserved
- return result != null
- && result.value().equals(value)
- && result.timestamp() == -1L
- && result.headers().equals(expectedHeaders);
- } catch (final Exception e) {
- e.printStackTrace();
- return false;
- }
- }, 60_000L, "Could not verify plain windowed value with headers in
time.");
+ },
+ "Could not verify plain window value with empty headers and
timestamp in time.");
}
private void processWindowedKeyValueAndVerifyTimestamped(final String key,
final String
value,
final long
timestamp) throws Exception {
- IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp(
- inputStream,
- singletonList(KeyValue.pair(key, value)),
- TestUtils.producerConfig(CLUSTER.bootstrapServers(),
- StringSerializer.class,
- StringSerializer.class),
- timestamp,
- false);
-
- TestUtils.waitForCondition(() -> {
- try {
- final ReadOnlyWindowStore<String, ValueAndTimestamp<String>>
store =
- IntegrationTestUtils.getStore(WINDOW_STORE_NAME,
kafkaStreams, QueryableStoreTypes.timestampedWindowStore());
-
- if (store == null) {
- return false;
- }
-
- final long windowStart = timestamp - (timestamp %
WINDOW_SIZE_MS);
- final ValueAndTimestamp<String> result = store.fetch(key,
windowStart);
+ produce(key, value, timestamp);
+ awaitStore(WINDOW_STORE_NAME, QueryableStoreTypes.<String,
String>timestampedWindowStore(),
+ store -> {
+ final ValueAndTimestamp<String> result = store.fetch(key,
windowStart(timestamp));
return result != null
&& result.value().equals(value)
&& result.timestamp() == timestamp;
- } catch (final Exception e) {
- return false;
- }
- }, 60_000L, "Could not verify timestamped value in time.");
+ },
+ "Could not verify timestamped value in time.");
}
private void processWindowedKeyValueWithHeadersAndVerify(final String key,
final String
value,
final long
timestamp,
final Headers
headers,
final Headers
expectedHeaders) throws Exception {
- IntegrationTestUtils.produceKeyValuesSynchronouslyWithTimestamp(
- inputStream,
- singletonList(KeyValue.pair(key, value)),
- TestUtils.producerConfig(CLUSTER.bootstrapServers(),
- StringSerializer.class,
- StringSerializer.class),
- headers,
- timestamp,
- false);
-
- TestUtils.waitForCondition(() -> {
- try {
- final ReadOnlyWindowStore<String,
ValueTimestampHeaders<String>> store =
- IntegrationTestUtils.getStore(WINDOW_STORE_NAME,
kafkaStreams, QueryableStoreTypes.timestampedWindowStoreWithHeaders());
-
- if (store == null) {
- return false;
- }
-
- final long windowStart = timestamp - (timestamp %
WINDOW_SIZE_MS);
-
- final List<KeyValue<Windowed<String>,
ValueTimestampHeaders<String>>> results = new LinkedList<>();
- try (final KeyValueIterator<Windowed<String>,
ValueTimestampHeaders<String>> iterator = store.all()) {
- while (iterator.hasNext()) {
- final KeyValue<Windowed<String>,
ValueTimestampHeaders<String>> kv = iterator.next();
- if (kv.key.key().equals(key) &&
kv.key.window().start() == windowStart) {
- results.add(kv);
- }
- }
- }
-
- if (results.isEmpty()) {
- return false;
- }
-
- final ValueTimestampHeaders<String> result =
results.get(0).value;
- return result != null
- && result.value().equals(value)
- && result.timestamp() == timestamp
- && result.headers().equals(expectedHeaders);
- } catch (final Exception e) {
- LOG.error("Error while verifying windowed value with headers",
e);
- return false;
- }
- }, 60_000L, "Could not verify windowed value with headers in time.");
+ processWindowedKeyValueWithHeadersAndVerify(key, value, timestamp,
timestamp, headers, expectedHeaders);
}
- private void verifyWindowValueWithEmptyHeaders(final String key,
- final String value,
- final long timestamp)
throws Exception {
- TestUtils.waitForCondition(() -> {
- try {
- final ReadOnlyWindowStore<String,
ValueTimestampHeaders<String>> store =
- IntegrationTestUtils.getStore(WINDOW_STORE_NAME,
kafkaStreams, QueryableStoreTypes.timestampedWindowStoreWithHeaders());
-
- if (store == null) {
- return false;
- }
-
- final long windowStart = timestamp - (timestamp %
WINDOW_SIZE_MS);
-
- final List<KeyValue<Windowed<String>,
ValueTimestampHeaders<String>>> results = new LinkedList<>();
- try (final KeyValueIterator<Windowed<String>,
ValueTimestampHeaders<String>> iterator = store.all()) {
- while (iterator.hasNext()) {
- final KeyValue<Windowed<String>,
ValueTimestampHeaders<String>> kv = iterator.next();
- if (kv.key.key().equals(key) &&
kv.key.window().start() == windowStart) {
- results.add(kv);
- }
- }
- }
-
- if (results.isEmpty()) {
- return false;
- }
-
- final ValueTimestampHeaders<String> result =
results.get(0).value;
- assertNotNull(result, "Result should not be null");
- assertEquals(value, result.value(), "Value should match");
- assertEquals(timestamp, result.timestamp(), "Timestamp should
match");
-
- // Verify headers exist but are empty (migrated from
timestamped store without headers)
- assertNotNull(result.headers(), "Headers should not be null
for migrated data");
- assertEquals(0, result.headers().toArray().length, "Headers
should be empty for migrated data");
-
- return true;
- } catch (final Exception e) {
- LOG.error("Error while verifying legacy value with empty
headers", e);
- return false;
- }
- }, 60_000L, "Could not verify legacy value with empty headers in
time.");
+ /**
+ * Produces a windowed record with headers and verifies the stored
value/headers, expecting
+ * {@code expectedTimestamp} in the store. For a plain window store (no
timestamp preserved)
+ * pass {@code expectedTimestamp == -1L}; otherwise pass the produced
{@code timestamp}.
+ */
+ private void processWindowedKeyValueWithHeadersAndVerify(final String key,
+ final String
value,
+ final long
timestamp,
+ final long
expectedTimestamp,
+ final Headers
headers,
+ final Headers
expectedHeaders) throws Exception {
+ produce(key, value, timestamp, headers);
+
+ awaitStore(WINDOW_STORE_NAME, QueryableStoreTypes.<String,
String>timestampedWindowStoreWithHeaders(),
Review Comment:
Missed it. Yea done for both now.
--
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]