|
31 | 31 | import org.apache.kafka.common.header.internals.RecordHeaders;
|
32 | 32 | import org.assertj.core.api.InstanceOfAssertFactories;
|
33 | 33 | import org.junit.jupiter.api.Test;
|
| 34 | +import org.junit.jupiter.params.ParameterizedTest; |
| 35 | +import org.junit.jupiter.params.provider.ValueSource; |
34 | 36 |
|
35 | 37 | import org.springframework.core.log.LogAccessor;
|
36 | 38 | import org.springframework.kafka.retrytopic.RetryTopicHeaders;
|
@@ -413,6 +415,103 @@ void multiValueHeaderToTest() {
|
413 | 415 | .containsExactly(multiValueWildCardHeader2Value1, multiValueWildCardHeader2Value2);
|
414 | 416 | }
|
415 | 417 |
|
| 418 | + @ParameterizedTest |
| 419 | + @ValueSource(ints = {2000}) |
| 420 | +// @ValueSource(ints = {500, 1000, 2000}) |
| 421 | + void hugeNumberOfSingleValueHeaderToTest(int numberOfSingleValueHeaderCount) { |
| 422 | + // GIVEN |
| 423 | + Headers rawHeaders = new RecordHeaders(); |
| 424 | + |
| 425 | + String multiValueHeader1 = "test-multi-value1"; |
| 426 | + byte[] multiValueHeader1Value1 = { 0, 0, 0, 0 }; |
| 427 | + byte[] multiValueHeader1Value2 = { 0, 0, 0, 1 }; |
| 428 | + |
| 429 | + rawHeaders.add(multiValueHeader1, multiValueHeader1Value1); |
| 430 | + rawHeaders.add(multiValueHeader1, multiValueHeader1Value2); |
| 431 | + |
| 432 | + byte[] deliveryAttemptHeaderValue = { 0, 0, 0, 1 }; |
| 433 | + byte[] originalOffsetHeaderValue = { 0, 0, 0, 2 }; |
| 434 | + byte[] defaultHeaderAttemptsValues = { 0, 0, 0, 5 }; |
| 435 | + |
| 436 | + rawHeaders.add(KafkaHeaders.DELIVERY_ATTEMPT, deliveryAttemptHeaderValue); |
| 437 | + rawHeaders.add(KafkaHeaders.ORIGINAL_OFFSET, originalOffsetHeaderValue); |
| 438 | + rawHeaders.add(RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS, defaultHeaderAttemptsValues); |
| 439 | + |
| 440 | + byte[] singleValueHeaderValue = { 0, 0, 0, 6 }; |
| 441 | + for (int i = 0; i < numberOfSingleValueHeaderCount; i++) { |
| 442 | + String singleValueHeader = "test-single-value" + i; |
| 443 | + rawHeaders.add(singleValueHeader, singleValueHeaderValue); |
| 444 | + } |
| 445 | + |
| 446 | + DefaultKafkaHeaderMapper mapper = new DefaultKafkaHeaderMapper(); |
| 447 | + mapper.setMultiValueHeaderPatterns(multiValueHeader1); |
| 448 | + |
| 449 | + // WHEN |
| 450 | + Map<String, Object> mappedHeaders = new HashMap<>(); |
| 451 | + mapper.toHeaders(rawHeaders, mappedHeaders); |
| 452 | + |
| 453 | + // THEN |
| 454 | + assertThat(mappedHeaders.get(KafkaHeaders.DELIVERY_ATTEMPT)).isEqualTo(1); |
| 455 | + assertThat(mappedHeaders.get(KafkaHeaders.ORIGINAL_OFFSET)).isEqualTo(originalOffsetHeaderValue); |
| 456 | + assertThat(mappedHeaders.get(RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS)).isEqualTo(defaultHeaderAttemptsValues); |
| 457 | + |
| 458 | + for (int i = 0; i < numberOfSingleValueHeaderCount; i++) { |
| 459 | + String singleValueHeader = "test-single-value" + i; |
| 460 | + assertThat(mappedHeaders.get(singleValueHeader)).isEqualTo(singleValueHeaderValue); |
| 461 | + } |
| 462 | + |
| 463 | + assertThat(mappedHeaders) |
| 464 | + .extractingByKey(multiValueHeader1, InstanceOfAssertFactories.list(byte[].class)) |
| 465 | + .containsExactly(multiValueHeader1Value1, multiValueHeader1Value2); |
| 466 | + } |
| 467 | + |
| 468 | + @ParameterizedTest |
| 469 | + @ValueSource(ints = {500, 1000, 2000}) |
| 470 | + void hugeNumberOfMultiValueHeaderToTest(int numberOfMultiValueHeaderCount) { |
| 471 | + // GIVEN |
| 472 | + DefaultKafkaHeaderMapper mapper = new DefaultKafkaHeaderMapper(); |
| 473 | + Headers rawHeaders = new RecordHeaders(); |
| 474 | + |
| 475 | + byte[] multiValueHeader1Value1 = { 0, 0, 0, 0 }; |
| 476 | + byte[] multiValueHeader1Value2 = { 0, 0, 0, 1 }; |
| 477 | + |
| 478 | + for (int i = 0; i < numberOfMultiValueHeaderCount; i++) { |
| 479 | + String multiValueHeader = "test-multi-value" + i; |
| 480 | + mapper.setMultiValueHeaderPatterns(multiValueHeader); |
| 481 | + rawHeaders.add(multiValueHeader, multiValueHeader1Value1); |
| 482 | + rawHeaders.add(multiValueHeader, multiValueHeader1Value2); |
| 483 | + } |
| 484 | + |
| 485 | + byte[] deliveryAttemptHeaderValue = { 0, 0, 0, 1 }; |
| 486 | + byte[] originalOffsetHeaderValue = { 0, 0, 0, 2 }; |
| 487 | + byte[] defaultHeaderAttemptsValues = { 0, 0, 0, 5 }; |
| 488 | + |
| 489 | + rawHeaders.add(KafkaHeaders.DELIVERY_ATTEMPT, deliveryAttemptHeaderValue); |
| 490 | + rawHeaders.add(KafkaHeaders.ORIGINAL_OFFSET, originalOffsetHeaderValue); |
| 491 | + rawHeaders.add(RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS, defaultHeaderAttemptsValues); |
| 492 | + |
| 493 | + String singleValueHeader = "test-single-value"; |
| 494 | + byte[] singleValueHeaderValue = { 0, 0, 0, 6 }; |
| 495 | + rawHeaders.add(singleValueHeader, singleValueHeaderValue); |
| 496 | + |
| 497 | + // WHEN |
| 498 | + Map<String, Object> mappedHeaders = new HashMap<>(); |
| 499 | + mapper.toHeaders(rawHeaders, mappedHeaders); |
| 500 | + |
| 501 | + // THEN |
| 502 | + assertThat(mappedHeaders.get(KafkaHeaders.DELIVERY_ATTEMPT)).isEqualTo(1); |
| 503 | + assertThat(mappedHeaders.get(KafkaHeaders.ORIGINAL_OFFSET)).isEqualTo(originalOffsetHeaderValue); |
| 504 | + assertThat(mappedHeaders.get(RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS)).isEqualTo(defaultHeaderAttemptsValues); |
| 505 | + assertThat(mappedHeaders.get(singleValueHeader)).isEqualTo(singleValueHeaderValue); |
| 506 | + |
| 507 | + for (int i = 0; i < numberOfMultiValueHeaderCount; i++) { |
| 508 | + String multiValueHeader = "test-multi-value" + i; |
| 509 | + assertThat(mappedHeaders) |
| 510 | + .extractingByKey(multiValueHeader, InstanceOfAssertFactories.list(byte[].class)) |
| 511 | + .containsExactly(multiValueHeader1Value1, multiValueHeader1Value2); |
| 512 | + } |
| 513 | + } |
| 514 | + |
416 | 515 | @Test
|
417 | 516 | void multiValueHeaderFromTest() {
|
418 | 517 | // GIVEN
|
|
0 commit comments