Skip to content

Commit 6007253

Browse files
committed
[Improve][Connector-V2] Migrate IoTDBv2 SQL dialect validation to OptionRule
1 parent 65e32aa commit 6007253

3 files changed

Lines changed: 99 additions & 38 deletions

File tree

seatunnel-connectors-v2/connector-iotdb-v2/src/main/java/org/apache/seatunnel/connectors/seatunnel/iotdbv2/sink/IoTDBv2SinkFactory.java

Lines changed: 8 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -18,15 +18,14 @@
1818
package org.apache.seatunnel.connectors.seatunnel.iotdbv2.sink;
1919

2020
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
21+
import org.apache.seatunnel.api.configuration.util.Conditions;
2122
import org.apache.seatunnel.api.configuration.util.OptionRule;
2223
import org.apache.seatunnel.api.table.connector.TableSink;
2324
import org.apache.seatunnel.api.table.factory.Factory;
2425
import org.apache.seatunnel.api.table.factory.TableSinkFactory;
2526
import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
26-
import org.apache.seatunnel.common.exception.CommonErrorCode;
2727
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.config.IoTDBv2SinkOptions;
2828
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.constant.SinkConstants;
29-
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.exception.IotdbConnectorException;
3029

3130
import com.google.auto.service.AutoService;
3231
import lombok.extern.slf4j.Slf4j;
@@ -50,6 +49,8 @@ public OptionRule optionRule() {
5049
IoTDBv2SinkOptions.KEY_DEVICE)
5150
.optional(
5251
IoTDBv2SinkOptions.SQL_DIALECT,
52+
Conditions.matches(IoTDBv2SinkOptions.SQL_DIALECT, "(?i)^(tree|table)$"))
53+
.optional(
5354
IoTDBv2SinkOptions.KEY_TIMESTAMP,
5455
IoTDBv2SinkOptions.KEY_TAG_FIELDS,
5556
IoTDBv2SinkOptions.KEY_ATTRIBUTE_FIELDS,
@@ -69,22 +70,11 @@ public OptionRule optionRule() {
6970
@Override
7071
public TableSink createSink(TableSinkFactoryContext context) {
7172
ReadonlyConfig conf = context.getOptions();
72-
String targetSqlDialect;
73-
if (conf.get(IoTDBv2SinkOptions.SQL_DIALECT) != null) {
74-
String sqlDialect = conf.get(IoTDBv2SinkOptions.SQL_DIALECT);
75-
if (SinkConstants.TABLE.equalsIgnoreCase(sqlDialect)) {
76-
targetSqlDialect = SinkConstants.TABLE;
77-
} else {
78-
if (SinkConstants.TREE.equalsIgnoreCase(sqlDialect)) {
79-
targetSqlDialect = SinkConstants.TREE;
80-
} else {
81-
throw new IotdbConnectorException(
82-
CommonErrorCode.ILLEGAL_ARGUMENT, "Sql dialect not supported");
83-
}
84-
}
85-
} else {
86-
targetSqlDialect = SinkConstants.TREE;
87-
}
73+
String sqlDialect = conf.get(IoTDBv2SinkOptions.SQL_DIALECT);
74+
String targetSqlDialect =
75+
SinkConstants.TABLE.equalsIgnoreCase(sqlDialect)
76+
? SinkConstants.TABLE
77+
: SinkConstants.TREE;
8878
return () ->
8979
new IoTDBv2Sink(context.getOptions(), context.getCatalogTable(), targetSqlDialect);
9080
}

seatunnel-connectors-v2/connector-iotdb-v2/src/main/java/org/apache/seatunnel/connectors/seatunnel/iotdbv2/source/IoTDBv2SourceFactory.java

Lines changed: 8 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
package org.apache.seatunnel.connectors.seatunnel.iotdbv2.source;
1919

2020
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
21+
import org.apache.seatunnel.api.configuration.util.Conditions;
2122
import org.apache.seatunnel.api.configuration.util.OptionRule;
2223
import org.apache.seatunnel.api.options.ConnectorCommonOptions;
2324
import org.apache.seatunnel.api.source.SeaTunnelSource;
@@ -28,10 +29,8 @@
2829
import org.apache.seatunnel.api.table.factory.Factory;
2930
import org.apache.seatunnel.api.table.factory.TableSourceFactory;
3031
import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
31-
import org.apache.seatunnel.common.exception.CommonErrorCode;
3232
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.config.IoTDBv2SourceOptions;
3333
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.constant.SourceConstants;
34-
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.exception.IotdbConnectorException;
3534

3635
import com.google.auto.service.AutoService;
3736

@@ -55,6 +54,8 @@ public OptionRule optionRule() {
5554
ConnectorCommonOptions.SCHEMA)
5655
.optional(
5756
IoTDBv2SourceOptions.SQL_DIALECT,
57+
Conditions.matches(IoTDBv2SourceOptions.SQL_DIALECT, "(?i)^(tree|table)$"))
58+
.optional(
5859
IoTDBv2SourceOptions.DATABASE,
5960
IoTDBv2SourceOptions.FETCH_SIZE,
6061
IoTDBv2SourceOptions.DEFAULT_THRIFT_BUFFER_SIZE,
@@ -71,22 +72,11 @@ public OptionRule optionRule() {
7172
TableSource<T, SplitT, StateT> createSource(TableSourceFactoryContext context) {
7273
CatalogTable catalogTable = CatalogTableUtil.buildWithConfig(context.getOptions());
7374
ReadonlyConfig conf = context.getOptions();
74-
String targetSqlDialect;
75-
if (conf.get(IoTDBv2SourceOptions.SQL_DIALECT) != null) {
76-
String sqlDialect = conf.get(IoTDBv2SourceOptions.SQL_DIALECT);
77-
if (SourceConstants.TABLE.equalsIgnoreCase(sqlDialect)) {
78-
targetSqlDialect = SourceConstants.TABLE;
79-
} else {
80-
if (SourceConstants.TREE.equalsIgnoreCase(sqlDialect)) {
81-
targetSqlDialect = SourceConstants.TREE;
82-
} else {
83-
throw new IotdbConnectorException(
84-
CommonErrorCode.ILLEGAL_ARGUMENT, "Sql dialect not supported");
85-
}
86-
}
87-
} else {
88-
targetSqlDialect = SourceConstants.TREE;
89-
}
75+
String sqlDialect = conf.get(IoTDBv2SourceOptions.SQL_DIALECT);
76+
String targetSqlDialect =
77+
SourceConstants.TABLE.equalsIgnoreCase(sqlDialect)
78+
? SourceConstants.TABLE
79+
: SourceConstants.TREE;
9080
return () ->
9181
(SeaTunnelSource<T, SplitT, StateT>)
9282
new IoTDBv2Source(catalogTable, context.getOptions(), targetSqlDialect);

seatunnel-connectors-v2/connector-iotdb-v2/src/test/java/org/apache/seatunnel/connectors/seatunnel/iotdbv2/IoTDBFactoryTest.java

Lines changed: 83 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,17 +17,98 @@
1717

1818
package org.apache.seatunnel.connectors.seatunnel.iotdbv2;
1919

20+
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
21+
import org.apache.seatunnel.api.configuration.util.ConfigValidator;
22+
import org.apache.seatunnel.api.configuration.util.OptionRule;
23+
import org.apache.seatunnel.api.configuration.util.OptionValidationException;
24+
import org.apache.seatunnel.api.options.ConnectorCommonOptions;
25+
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.config.IoTDBv2CommonOptions;
26+
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.config.IoTDBv2SinkOptions;
27+
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.config.IoTDBv2SourceOptions;
2028
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.sink.IoTDBv2SinkFactory;
2129
import org.apache.seatunnel.connectors.seatunnel.iotdbv2.source.IoTDBv2SourceFactory;
2230

2331
import org.junit.jupiter.api.Assertions;
2432
import org.junit.jupiter.api.Test;
2533

34+
import java.util.Arrays;
35+
import java.util.Collections;
36+
import java.util.HashMap;
37+
import java.util.Map;
38+
2639
class IoTDBFactoryTest {
2740

41+
private final OptionRule sourceRule = new IoTDBv2SourceFactory().optionRule();
42+
private final OptionRule sinkRule = new IoTDBv2SinkFactory().optionRule();
43+
2844
@Test
2945
void optionRule() {
30-
Assertions.assertNotNull((new IoTDBv2SourceFactory()).optionRule());
31-
Assertions.assertNotNull((new IoTDBv2SinkFactory()).optionRule());
46+
Assertions.assertNotNull(sourceRule);
47+
Assertions.assertNotNull(sinkRule);
48+
}
49+
50+
@Test
51+
void omittedSqlDialectIsAccepted() {
52+
Assertions.assertDoesNotThrow(() -> validate(validSourceConfig(), sourceRule));
53+
Assertions.assertDoesNotThrow(() -> validate(validSinkConfig(), sinkRule));
54+
}
55+
56+
@Test
57+
void supportedSqlDialectsAreAcceptedCaseInsensitively() {
58+
for (String dialect : Arrays.asList("tree", "table", "TrEe", "TaBlE")) {
59+
Map<String, Object> sourceConfig = validSourceConfig();
60+
sourceConfig.put(IoTDBv2CommonOptions.SQL_DIALECT.key(), dialect);
61+
Assertions.assertDoesNotThrow(() -> validate(sourceConfig, sourceRule));
62+
63+
Map<String, Object> sinkConfig = validSinkConfig();
64+
sinkConfig.put(IoTDBv2CommonOptions.SQL_DIALECT.key(), dialect);
65+
Assertions.assertDoesNotThrow(() -> validate(sinkConfig, sinkRule));
66+
}
67+
}
68+
69+
@Test
70+
void unsupportedSqlDialectIsRejected() {
71+
Map<String, Object> sourceConfig = validSourceConfig();
72+
sourceConfig.put(IoTDBv2CommonOptions.SQL_DIALECT.key(), "unsupported");
73+
OptionValidationException sourceException =
74+
Assertions.assertThrows(
75+
OptionValidationException.class, () -> validate(sourceConfig, sourceRule));
76+
Assertions.assertTrue(
77+
sourceException.getMessage().contains(IoTDBv2CommonOptions.SQL_DIALECT.key()));
78+
79+
Map<String, Object> sinkConfig = validSinkConfig();
80+
sinkConfig.put(IoTDBv2CommonOptions.SQL_DIALECT.key(), "unsupported");
81+
OptionValidationException sinkException =
82+
Assertions.assertThrows(
83+
OptionValidationException.class, () -> validate(sinkConfig, sinkRule));
84+
Assertions.assertTrue(
85+
sinkException.getMessage().contains(IoTDBv2CommonOptions.SQL_DIALECT.key()));
86+
}
87+
88+
private void validate(Map<String, Object> config, OptionRule rule) {
89+
ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(rule);
90+
}
91+
92+
private Map<String, Object> validSourceConfig() {
93+
Map<String, Object> config = commonConfig();
94+
config.put(IoTDBv2SourceOptions.SQL.key(), "select * from root.test");
95+
config.put(ConnectorCommonOptions.SCHEMA.key(), Collections.emptyMap());
96+
return config;
97+
}
98+
99+
private Map<String, Object> validSinkConfig() {
100+
Map<String, Object> config = commonConfig();
101+
config.put(IoTDBv2SinkOptions.STORAGE_GROUP.key(), "root.test");
102+
config.put(IoTDBv2SinkOptions.KEY_DEVICE.key(), "device");
103+
return config;
104+
}
105+
106+
private Map<String, Object> commonConfig() {
107+
Map<String, Object> config = new HashMap<>();
108+
config.put(
109+
IoTDBv2CommonOptions.NODE_URLS.key(), Collections.singletonList("127.0.0.1:6667"));
110+
config.put(IoTDBv2CommonOptions.USERNAME.key(), "root");
111+
config.put(IoTDBv2CommonOptions.PASSWORD.key(), "root");
112+
return config;
32113
}
33114
}

0 commit comments

Comments
 (0)