Skip to content

Commit b5c468f

Browse files
KAFKA-18115; Fix for loading big files while performing load tests (#18391)
When performing perf tests, we can specify a payload using the "--payloadFile" flag. This file is utilized during the load/performance testing process. This causes the entire file to get loaded into a String and split using the delimiter. However, if the file is large, it may result in NegativeArraySizeException error. Moving the file loading logic to Scanner which doesn't have this issue. Reviewers: José Armando García Sancio <jsancio@apache.org>, Ken Huang <s7133700@gmail.com>, Zhe Guang <zheguang.zhao@alumni.brown.edu>
1 parent 99ecd5c commit b5c468f

2 files changed

Lines changed: 85 additions & 5 deletions

File tree

tools/src/main/java/org/apache/kafka/tools/ProducerPerformance.java

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@
4242
import java.util.List;
4343
import java.util.Optional;
4444
import java.util.Properties;
45+
import java.util.Scanner;
4546
import java.util.SplittableRandom;
4647

4748
import static net.sourceforge.argparse4j.impl.Arguments.store;
@@ -194,13 +195,17 @@ static List<byte[]> readPayloadFile(String payloadFilePath, String payloadDelimi
194195
throw new IllegalArgumentException("File does not exist or empty file provided.");
195196
}
196197

197-
String[] payloadList = Files.readString(path).split(payloadDelimiter);
198+
try (Scanner payLoadScanner = new Scanner(path, StandardCharsets.UTF_8)) {
199+
//setting the delimiter while parsing the file, avoids loading entire data in memory before split
200+
payLoadScanner.useDelimiter(payloadDelimiter);
201+
while (payLoadScanner.hasNext()) {
202+
byte[] payloadBytes = payLoadScanner.next().getBytes(StandardCharsets.UTF_8);
203+
payloadByteList.add(payloadBytes);
204+
}
205+
}
198206

199-
System.out.println("Number of messages read: " + payloadList.length);
207+
System.out.println("Number of messages read: " + payloadByteList.size());
200208

201-
for (String payload : payloadList) {
202-
payloadByteList.add(payload.getBytes(StandardCharsets.UTF_8));
203-
}
204209
}
205210
return payloadByteList;
206211
}

tools/src/test/java/org/apache/kafka/tools/ProducerPerformanceTest.java

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,81 @@ public void testReadProps() throws Exception {
104104
Utils.delete(producerConfig);
105105
}
106106

107+
@Test
108+
public void testReadPayloadFileWithAlternateDelimiters() throws Exception {
109+
List<byte[]> payloadByteList;
110+
111+
payloadByteList = generateListFromFileUsingDelimiter("Hello~~Kafka", "~~");
112+
assertEquals(2, payloadByteList.size());
113+
assertEquals("Hello", new String(payloadByteList.get(0)));
114+
assertEquals("Kafka", new String(payloadByteList.get(1)));
115+
116+
payloadByteList = generateListFromFileUsingDelimiter("Hello,Kafka,", ",");
117+
assertEquals(2, payloadByteList.size());
118+
assertEquals("Hello", new String(payloadByteList.get(0)));
119+
assertEquals("Kafka", new String(payloadByteList.get(1)));
120+
121+
payloadByteList = generateListFromFileUsingDelimiter("Hello\t\tKafka", "\t");
122+
assertEquals(3, payloadByteList.size());
123+
assertEquals("Hello", new String(payloadByteList.get(0)));
124+
assertEquals("Kafka", new String(payloadByteList.get(2)));
125+
126+
payloadByteList = generateListFromFileUsingDelimiter("Hello\n\nKafka\n", "\n");
127+
assertEquals(3, payloadByteList.size());
128+
assertEquals("Hello", new String(payloadByteList.get(0)));
129+
assertEquals("Kafka", new String(payloadByteList.get(2)));
130+
131+
payloadByteList = generateListFromFileUsingDelimiter("Hello::Kafka::World", "\\s*::\\s*");
132+
assertEquals(3, payloadByteList.size());
133+
assertEquals("Hello", new String(payloadByteList.get(0)));
134+
assertEquals("Kafka", new String(payloadByteList.get(1)));
135+
136+
}
137+
138+
@Test
139+
public void testCompareStringSplitWithScannerDelimiter() throws Exception {
140+
141+
String contents = "Hello~~Kafka";
142+
String payloadDelimiter = "~~";
143+
compareList(generateListFromFileUsingDelimiter(contents, payloadDelimiter), contents.split(payloadDelimiter));
144+
145+
contents = "Hello,Kafka,";
146+
payloadDelimiter = ",";
147+
compareList(generateListFromFileUsingDelimiter(contents, payloadDelimiter), contents.split(payloadDelimiter));
148+
149+
contents = "Hello\t\tKafka";
150+
payloadDelimiter = "\t";
151+
compareList(generateListFromFileUsingDelimiter(contents, payloadDelimiter), contents.split(payloadDelimiter));
152+
153+
contents = "Hello\n\nKafka\n";
154+
payloadDelimiter = "\n";
155+
compareList(generateListFromFileUsingDelimiter(contents, payloadDelimiter), contents.split(payloadDelimiter));
156+
157+
contents = "Hello::Kafka::World";
158+
payloadDelimiter = "\\s*::\\s*";
159+
compareList(generateListFromFileUsingDelimiter(contents, payloadDelimiter), contents.split(payloadDelimiter));
160+
161+
}
162+
163+
private void compareList(List<byte[]> payloadByteList, String[] payloadByteListFromSplit) {
164+
assertEquals(payloadByteListFromSplit.length, payloadByteList.size());
165+
for (int i = 0; i < payloadByteListFromSplit.length; i++) {
166+
assertEquals(payloadByteListFromSplit[i], new String(payloadByteList.get(i)));
167+
}
168+
}
169+
170+
private List<byte[]> generateListFromFileUsingDelimiter(String fileContent, String payloadDelimiter) throws Exception {
171+
File payloadFile = null;
172+
List<byte[]> payloadByteList;
173+
try {
174+
payloadFile = createTempFile(fileContent);
175+
payloadByteList = ProducerPerformance.readPayloadFile(payloadFile.getAbsolutePath(), payloadDelimiter);
176+
} finally {
177+
Utils.delete(payloadFile);
178+
}
179+
return payloadByteList;
180+
}
181+
107182
@Test
108183
public void testNumberOfCallsForSendAndClose() throws IOException {
109184
doReturn(null).when(producerMock).send(any(), any());

0 commit comments

Comments
 (0)