Skip to content

Commit a875c19

Browse files
authored
add jsonl as a supported extension for S3 sink (opensearch-project#6159)
* add jsonl as a supported extension for S3 sink Signed-off-by: Xun Zhang <xunzh@amazon.com> * make ExtensionOption an enum Signed-off-by: Xun Zhang <xunzh@amazon.com> * address more comments Signed-off-by: Xun Zhang <xunzh@amazon.com> --------- Signed-off-by: Xun Zhang <xunzh@amazon.com>
1 parent 82f56d2 commit a875c19

5 files changed

Lines changed: 148 additions & 3 deletions

File tree

Original file line numberDiff line numberDiff line change
@@ -0,0 +1,40 @@
1+
/*
2+
* Copyright OpenSearch Contributors
3+
* SPDX-License-Identifier: Apache-2.0
4+
*/
5+
6+
package org.opensearch.dataprepper.plugins.codec.json;
7+
8+
import com.fasterxml.jackson.annotation.JsonCreator;
9+
import com.fasterxml.jackson.annotation.JsonValue;
10+
11+
import java.util.Arrays;
12+
import java.util.Map;
13+
import java.util.stream.Collectors;
14+
15+
public enum ExtensionOption {
16+
NDJSON("ndjson"),
17+
JSONL("jsonl");
18+
19+
private final String extension;
20+
21+
private static final Map<String, ExtensionOption> EXTENSIONS_MAP = Arrays.stream(ExtensionOption.values())
22+
.collect(Collectors.toMap(
23+
ExtensionOption::getExtension,
24+
extensionOption -> extensionOption
25+
));
26+
27+
ExtensionOption(String extension) {
28+
this.extension = extension;
29+
}
30+
31+
@JsonValue
32+
public String getExtension() {
33+
return extension;
34+
}
35+
36+
@JsonCreator
37+
public static ExtensionOption fromExtension(String extension) {
38+
return EXTENSIONS_MAP.get(extension);
39+
}
40+
}

data-prepper-plugins/parse-json-processor/src/main/java/org/opensearch/dataprepper/plugins/codec/json/NdjsonOutputCodec.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,12 +20,13 @@
2020
*/
2121
@DataPrepperPlugin(name = "ndjson", pluginType = OutputCodec.class, pluginConfigurationType = NdjsonOutputConfig.class)
2222
public class NdjsonOutputCodec implements OutputCodec {
23-
private static final String NDJSON = "ndjson";
23+
private final NdjsonOutputConfig ndjsonOutputConfig;
2424
private OutputCodecContext deprecatedSupportCodecContext;
2525

2626
@DataPrepperPluginConstructor
2727
public NdjsonOutputCodec(final NdjsonOutputConfig config) {
28-
Objects.requireNonNull(config);
28+
Objects.requireNonNull(config, "NdjsonOutputConfig cannot be null");
29+
this.ndjsonOutputConfig = config;
2930
}
3031

3132
private static class NdjsonWriter implements Writer {
@@ -83,7 +84,7 @@ public void complete(final OutputStream outputStream) throws IOException {
8384

8485
@Override
8586
public String getExtension() {
86-
return NDJSON;
87+
return ndjsonOutputConfig.getExtensionOption().getExtension();
8788
}
8889

8990
private static void doWriteEvent(final OutputStream outputStream, final Event event, final OutputCodecContext codecContext) throws IOException {

data-prepper-plugins/parse-json-processor/src/main/java/org/opensearch/dataprepper/plugins/codec/json/NdjsonOutputConfig.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,18 @@
55

66
package org.opensearch.dataprepper.plugins.codec.json;
77

8+
import com.fasterxml.jackson.annotation.JsonProperty;
9+
import com.fasterxml.jackson.annotation.JsonPropertyDescription;
10+
811
/**
912
* Configuration class for the newline delimited codec.
1013
*/
1114
public class NdjsonOutputConfig {
15+
@JsonProperty("extension")
16+
@JsonPropertyDescription("Defines the file extension of the file produced by the sink. Default is 'ndjson'.")
17+
private ExtensionOption extension = ExtensionOption.NDJSON;
18+
19+
public ExtensionOption getExtensionOption() {
20+
return extension != null ? extension : ExtensionOption.NDJSON;
21+
}
1222
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
/*
2+
* Copyright OpenSearch Contributors
3+
* SPDX-License-Identifier: Apache-2.0
4+
*/
5+
6+
package org.opensearch.dataprepper.plugins.codec.json;
7+
8+
import org.junit.jupiter.api.extension.ExtensionContext;
9+
import org.junit.jupiter.params.ParameterizedTest;
10+
import org.junit.jupiter.params.provider.Arguments;
11+
import org.junit.jupiter.params.provider.ArgumentsProvider;
12+
import org.junit.jupiter.params.provider.ArgumentsSource;
13+
import org.junit.jupiter.params.provider.EnumSource;
14+
15+
import java.util.stream.Stream;
16+
17+
import static org.hamcrest.CoreMatchers.equalTo;
18+
import static org.hamcrest.CoreMatchers.not;
19+
import static org.hamcrest.CoreMatchers.notNullValue;
20+
import static org.hamcrest.MatcherAssert.assertThat;
21+
import static org.hamcrest.Matchers.emptyString;
22+
import static org.junit.jupiter.params.provider.Arguments.arguments;
23+
24+
public class ExtensionOptionTest {
25+
@ParameterizedTest
26+
@EnumSource(ExtensionOption.class)
27+
void fromExtension_returns_expected_value(final ExtensionOption extensionOption) {
28+
assertThat(ExtensionOption.fromExtension(extensionOption.getExtension()), equalTo(extensionOption));
29+
}
30+
31+
@ParameterizedTest
32+
@EnumSource(ExtensionOption.class)
33+
void getExtension_returns_non_empty_null_for_all_types(final ExtensionOption extensionOption) {
34+
assertThat(extensionOption.getExtension(), notNullValue());
35+
}
36+
37+
@ParameterizedTest
38+
@EnumSource(value = ExtensionOption.class)
39+
void getExtension_returns_non_empty_string_for_all_types_except_none(final ExtensionOption extensionOption) {
40+
assertThat(extensionOption.getExtension(), notNullValue());
41+
assertThat(extensionOption.getExtension(), not(emptyString()));
42+
}
43+
44+
@ParameterizedTest
45+
@ArgumentsSource(ExtensionOptionToKnownName.class)
46+
void getTransformName_returns_expected_name(final ExtensionOption extensionOption, final String expectedString) {
47+
assertThat(extensionOption.getExtension(), equalTo(expectedString));
48+
}
49+
50+
static class ExtensionOptionToKnownName implements ArgumentsProvider {
51+
@Override
52+
public Stream<? extends Arguments> provideArguments(final ExtensionContext extensionContext) {
53+
return Stream.of(
54+
arguments(ExtensionOption.NDJSON, "ndjson"),
55+
arguments(ExtensionOption.JSONL, "jsonl")
56+
);
57+
}
58+
}
59+
}

data-prepper-plugins/parse-json-processor/src/test/java/org/opensearch/dataprepper/plugins/codec/json/NdjsonOutputCodecTest.java

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,10 @@
4141
import static org.hamcrest.CoreMatchers.notNullValue;
4242
import static org.hamcrest.MatcherAssert.assertThat;
4343
import static org.hamcrest.Matchers.anEmptyMap;
44+
import static org.junit.jupiter.api.Assertions.assertThrows;
4445
import static org.mockito.Mockito.mock;
4546
import static org.mockito.Mockito.verifyNoInteractions;
47+
import static org.mockito.Mockito.when;
4648

4749
@ExtendWith(MockitoExtension.class)
4850
class NdjsonOutputCodecTest {
@@ -121,6 +123,39 @@ void write_single_using_writer() throws IOException {
121123
assertThat(serializedMap, equalTo(eventMap));
122124
}
123125

126+
@Test
127+
void constructor_with_valid_ndjson_extension_should_not_throw_exception() {
128+
when(config.getExtensionOption()).thenReturn(ExtensionOption.NDJSON);
129+
130+
final NdjsonOutputCodec codec = createObjectUnderTest();
131+
132+
assertThat(codec.getExtension(), equalTo("ndjson"));
133+
}
134+
135+
@Test
136+
void constructor_with_valid_jsonl_extension_should_not_throw_exception() {
137+
when(config.getExtensionOption()).thenReturn(ExtensionOption.JSONL);
138+
139+
final NdjsonOutputCodec codec = createObjectUnderTest();
140+
141+
assertThat(codec.getExtension(), equalTo("jsonl"));
142+
}
143+
144+
@Test
145+
void constructor_with_uppercase_valid_extension_should_not_throw_exception() {
146+
when(config.getExtensionOption()).thenReturn(ExtensionOption.NDJSON);
147+
148+
final NdjsonOutputCodec codec = createObjectUnderTest();
149+
150+
assertThat(codec.getExtension(), equalTo("ndjson"));
151+
}
152+
153+
@Test
154+
void constructor_with_null_config_should_throw_exception() {
155+
assertThrows(NullPointerException.class, () -> new NdjsonOutputCodec(null),
156+
"NdjsonOutputConfig cannot be null");
157+
}
158+
124159
@ParameterizedTest
125160
@ValueSource(ints = {2, 100})
126161
void write_multiple(final int numberOfEvents) throws IOException {

0 commit comments

Comments
 (0)