From e37d0ca876856a74174e4da8146c02efb91d9f70 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Wed, 10 Jul 2019 15:25:28 -0700 Subject: [PATCH 1/2] KAFKA-8570: Grow buffer to hold down converted records if it was insufficiently sized --- .../kafka/common/record/AbstractRecords.java | 1 + .../kafka/common/record/FileRecordsTest.java | 36 +++++++++++++++++++ 2 files changed, 37 insertions(+) diff --git a/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java b/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java index 89a5413e00cf4..56c1cdbc85156 100644 --- a/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java +++ b/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java @@ -104,6 +104,7 @@ protected ConvertedRecords downConvert(Iterable offsets = asList(0L, 1L); + List magic = asList(RecordBatch.MAGIC_VALUE_V2, RecordBatch.MAGIC_VALUE_V1); // downgrade message format from v2 to v1 + List records = asList( + new SimpleRecord(1L, "k1".getBytes(), bytes), + new SimpleRecord(2L, "k2".getBytes(), bytes)); + byte toMagic = 1; + + // create MemoryRecords + ByteBuffer buffer = ByteBuffer.allocate(8000); + for (int i = 0; i < records.size(); i++) { + MemoryRecordsBuilder builder = MemoryRecords.builder(buffer, magic.get(i), compressionType, TimestampType.CREATE_TIME, 0L); + builder.appendWithOffset(offsets.get(i), records.get(i)); + builder.close(); + } + buffer.flip(); + + // create FileRecords, down-convert and verify + try (FileRecords fileRecords = FileRecords.open(tempFile())) { + fileRecords.append(MemoryRecords.readableRecords(buffer)); + fileRecords.flush(); + + Records convertedRecords = fileRecords.downConvert(toMagic, 0, time).records(); + verifyConvertedRecords(records, offsets, convertedRecords, compressionType, toMagic); + } + } + private void doTestConversion(CompressionType compressionType, byte toMagic) throws IOException { List offsets = asList(0L, 2L, 3L, 9L, 11L, 15L, 16L, 17L, 22L, 24L); From c41c216b6f42927d4de025203f3675af9c7dae7b Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Wed, 10 Jul 2019 15:27:41 -0700 Subject: [PATCH 2/2] Uncomment test change --- .../java/org/apache/kafka/common/record/AbstractRecords.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java b/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java index 56c1cdbc85156..0552e6b347dda 100644 --- a/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java +++ b/clients/src/main/java/org/apache/kafka/common/record/AbstractRecords.java @@ -104,7 +104,7 @@ protected ConvertedRecords downConvert(Iterable