This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new d9dd907735 [INLONG-8183][Agent] Optimize agent UT (#8185)
d9dd907735 is described below
commit d9dd9077357f839cb449d9045cf6d77cc1e3315b
Author: doleyzi <[email protected]>
AuthorDate: Wed Jun 7 04:36:59 2023 -0700
[INLONG-8183][Agent] Optimize agent UT (#8185)
---
.../apache/inlong/agent/task/TestTaskWrapper.java | 2 +-
.../agent/plugin/sources/TestTextFileReader.java | 13 +++----
.../inlong/agent/plugin/task/TestTextFileTask.java | 42 ++++++++++++----------
3 files changed, 31 insertions(+), 26 deletions(-)
diff --git
a/inlong-agent/agent-core/src/test/java/org/apache/inlong/agent/task/TestTaskWrapper.java
b/inlong-agent/agent-core/src/test/java/org/apache/inlong/agent/task/TestTaskWrapper.java
index e8433c4194..74f1c435c4 100755
---
a/inlong-agent/agent-core/src/test/java/org/apache/inlong/agent/task/TestTaskWrapper.java
+++
b/inlong-agent/agent-core/src/test/java/org/apache/inlong/agent/task/TestTaskWrapper.java
@@ -124,7 +124,7 @@ public class TestTaskWrapper {
@Override
public boolean isFinished() {
- return count > 10;
+ return count > 2;
}
@Override
diff --git
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sources/TestTextFileReader.java
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sources/TestTextFileReader.java
index 715f8d4c79..3979388358 100755
---
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sources/TestTextFileReader.java
+++
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sources/TestTextFileReader.java
@@ -244,12 +244,12 @@ public class TestTextFileReader {
Path localPath = Paths.get(testDir.toString(), "test.txt");
LOGGER.info("start to create {}", localPath);
List<String> beforeList = new ArrayList<>();
- for (int i = 0; i < 10; i++) {
+ for (int i = 0; i < 3; i++) {
beforeList.add("world");
}
Files.write(localPath, beforeList, StandardOpenOption.CREATE);
List<String> afterList = new ArrayList<>();
- for (int i = 0; i < 10; i++) {
+ for (int i = 0; i < 3; i++) {
afterList.add("world");
}
Files.write(localPath, afterList, StandardOpenOption.APPEND);
@@ -282,7 +282,7 @@ public class TestTextFileReader {
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
try {
List<String> beforeList = new ArrayList<>();
- for (int i = 0; i < 1000; i++) {
+ for (int i = 0; i < 3; i++) {
beforeList.add("hello, this is a new line for
testTextSeekReader");
}
Files.write(localPath, beforeList, StandardOpenOption.CREATE,
StandardOpenOption.APPEND);
@@ -290,12 +290,13 @@ public class TestTextFileReader {
LOGGER.info("ignored Exception ", ignored);
}
});
- TimeUnit.SECONDS.sleep(5);
+ TimeUnit.SECONDS.sleep(3);
int count = 0;
- while (!reader.isFinished() && count < 1000) {
+ while (!reader.isFinished() && count < 5) {
count += 1;
+ LOGGER.info("ignored count ", count);
}
- Assert.assertEquals(1000, count);
+ Assert.assertEquals(5, count);
}
private String getContent(String message) {
diff --git
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/task/TestTextFileTask.java
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/task/TestTextFileTask.java
index 4353c8de87..c0fec85586 100644
---
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/task/TestTextFileTask.java
+++
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/task/TestTextFileTask.java
@@ -170,7 +170,7 @@ public class TestTextFileTask {
public void testReadFull() throws IOException {
File file = TMP_FOLDER.newFile();
StringBuffer sb = new StringBuffer();
- String testData1 = IntStream.range(0, 100)
+ String testData1 = IntStream.range(0, 5)
.mapToObj(String::valueOf)
.collect(Collectors.joining(System.lineSeparator()));
sb.append(testData1);
@@ -186,21 +186,23 @@ public class TestTextFileTask {
jobProfile.set(JOB_FILE_META_ENV_LIST, ENV_CVM);
// mock data
final MockSink sink = mockTextTask(jobProfile);
- await().atMost(10, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 100);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 5);
await().atMost(10, TimeUnit.SECONDS).until(() ->
MonitorTextFile.getInstance().monitorNum() == 1);
- String testData = IntStream.range(100, 300)
+ String testData = IntStream.range(5, 10)
.mapToObj(String::valueOf)
.collect(Collectors.joining(System.lineSeparator()));
sb.append(testData);
sb.append(System.lineSeparator());
TestUtils.write(file.getAbsolutePath(), sb);
- await().atMost(10, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 100);
- String collectData = sink.getResult().stream().map(message -> {
- String content = new String(message.getBody(),
StandardCharsets.UTF_8);
- Map<String, String> logJson = GSON.fromJson(content, Map.class);
- return logJson.get(MetadataConstants.DATA_CONTENT);
- }).collect(Collectors.joining(System.lineSeparator()));
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 5);
+ synchronized (this) {
+ String collectData = sink.getResult().stream().map(message -> {
+ String content = new String(message.getBody(),
StandardCharsets.UTF_8);
+ Map<String, String> logJson = GSON.fromJson(content,
Map.class);
+ return logJson.get(MetadataConstants.DATA_CONTENT);
+ }).collect(Collectors.joining(System.lineSeparator()));
+ }
}
/**
@@ -210,7 +212,7 @@ public class TestTextFileTask {
public void testReadIncrement() throws IOException {
File file = TMP_FOLDER.newFile();
StringBuffer sb = new StringBuffer();
- sb.append(IntStream.range(0, 100)
+ sb.append(IntStream.range(0, 5)
.mapToObj(String::valueOf)
.collect(Collectors.joining(System.lineSeparator())));
sb.append(System.lineSeparator());
@@ -227,26 +229,28 @@ public class TestTextFileTask {
// mock data
final MockSink sink = mockTextTask(jobProfile);
await().atMost(10, TimeUnit.SECONDS).until(() ->
MonitorTextFile.getInstance().monitorNum() == 1);
- String testData = IntStream.range(100, 300)
+ String testData = IntStream.range(5, 10)
.mapToObj(String::valueOf)
.collect(Collectors.joining(System.lineSeparator()));
sb.append(testData);
sb.append(System.lineSeparator());
TestUtils.write(file.getAbsolutePath(), sb);
- await().atMost(10, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 300);
- String collectData = sink.getResult().stream().map(message -> {
- String content = new String(message.getBody(),
StandardCharsets.UTF_8);
- Map<String, String> logJson = GSON.fromJson(content, Map.class);
- return logJson.get(MetadataConstants.DATA_CONTENT);
- }).collect(Collectors.joining(System.lineSeparator()));
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 10);
+ synchronized (this) {
+ String collectData = sink.getResult().stream().map(message -> {
+ String content = new String(message.getBody(),
StandardCharsets.UTF_8);
+ Map<String, String> logJson = GSON.fromJson(content,
Map.class);
+ return logJson.get(MetadataConstants.DATA_CONTENT);
+ }).collect(Collectors.joining(System.lineSeparator()));
+ }
}
@Test
public void testScaleData() throws IOException {
File file = TMP_FOLDER.newFile();
StringBuffer sb = new StringBuffer();
- String testData1 = IntStream.range(0, 100)
+ String testData1 = IntStream.range(0, 5)
.mapToObj(String::valueOf)
.collect(Collectors.joining(System.lineSeparator()));
sb.append(testData1);
@@ -260,6 +264,6 @@ public class TestTextFileTask {
jobProfile.set(JobConstants.JOB_FILE_CONTENT_COLLECT_TYPE,
DataCollectType.FULL);
// mock data
final MockSink sink = mockTextTask(jobProfile);
- await().atMost(100, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 100);
+ await().atMost(100, TimeUnit.SECONDS).until(() ->
sink.getResult().size() == 5);
}
}