Skip to content

Commit 993f751

Browse files
committed
expose table option and client option.
1 parent 3d73038 commit 993f751

9 files changed

Lines changed: 599 additions & 271 deletions

File tree

flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/pom.xml

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,28 @@ limitations under the License.
9898
<version>${flink.version}</version>
9999
<scope>test</scope>
100100
</dependency>
101+
<dependency>
102+
<groupId>org.apache.flink</groupId>
103+
<artifactId> flink-cdc-composer</artifactId>
104+
<version>${project.version}</version>
105+
<scope>test</scope>
106+
</dependency>
107+
<dependency>
108+
<groupId>org.apache.flink</groupId>
109+
<artifactId> flink-cdc-composer</artifactId>
110+
<version>${project.version}</version>
111+
<type>test-jar</type>
112+
<scope>test</scope>
113+
</dependency>
114+
<dependency>
115+
<groupId>org.apache.flink</groupId>
116+
<artifactId>flink-cdc-pipeline-connector-values</artifactId>
117+
<version>${project.version}</version>
118+
<type>test-jar</type>
119+
<scope>test</scope>
120+
</dependency>
121+
122+
101123

102124
</dependencies>
103125

flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/factory/FlussDataSinkFactory.java

Lines changed: 35 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import org.apache.flink.cdc.common.sink.DataSink;
2424
import org.apache.flink.cdc.connectors.fluss.sink.FlussDataSink;
2525

26+
import com.alibaba.fluss.config.ConfigOptions;
2627
import com.alibaba.fluss.config.Configuration;
2728

2829
import java.util.HashMap;
@@ -36,37 +37,27 @@
3637
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.BOOTSTRAP_SERVERS;
3738
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.BUCKET_KEY;
3839
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.BUCKET_NUMBER;
39-
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.CLIENT_ID;
40-
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.PREFIX_FLUSS_PROPERTIES;
40+
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.CLIENT_PROPERTIES_PREFIX;
41+
import static org.apache.flink.cdc.connectors.fluss.sink.FlussDataSinkOptions.TABLE_PROPERTIES_PREFIX;
4142

4243
/** Factory for creating configured instances of {@link FlussDataSink}. */
4344
public class FlussDataSinkFactory implements DataSinkFactory {
4445
public static final String IDENTIFIER = "fluss";
4546

4647
@Override
4748
public DataSink createDataSink(Context context) {
48-
FactoryHelper.createFactoryHelper(this, context).validateExcept(PREFIX_FLUSS_PROPERTIES);
49+
FactoryHelper.createFactoryHelper(this, context)
50+
.validateExcept(CLIENT_PROPERTIES_PREFIX, TABLE_PROPERTIES_PREFIX);
4951
org.apache.flink.cdc.common.configuration.Configuration factoryConfiguration =
5052
context.getFactoryConfiguration();
5153

52-
Map<String, String> flussOptions = new HashMap<>();
53-
factoryConfiguration
54-
.toMap()
55-
.forEach(
56-
(key, value) -> {
57-
if (key.startsWith(PREFIX_FLUSS_PROPERTIES)) {
58-
flussOptions.put(
59-
key.substring(PREFIX_FLUSS_PROPERTIES.length()), value);
60-
} else {
61-
flussOptions.put(key, value);
62-
}
63-
});
64-
Configuration config = Configuration.fromMap(flussOptions);
54+
Configuration flussClientConfig = toFlussClientConfig(factoryConfiguration.toMap());
55+
Map<String, String> tableProperties = toFlussTableOptions(factoryConfiguration.toMap());
6556
Map<String, List<String>> bucketKeysMap =
6657
parseBucketKeys(factoryConfiguration.get(BUCKET_KEY));
6758
Map<String, Integer> bucketNumMap =
6859
parseBucketNumber(factoryConfiguration.get(BUCKET_NUMBER));
69-
return new FlussDataSink(config, bucketKeysMap, bucketNumMap);
60+
return new FlussDataSink(flussClientConfig, tableProperties, bucketKeysMap, bucketNumMap);
7061
}
7162

7263
@Override
@@ -86,7 +77,33 @@ public Set<ConfigOption<?>> optionalOptions() {
8677
Set<ConfigOption<?>> options = new HashSet<>();
8778
options.add(BUCKET_KEY);
8879
options.add(BUCKET_NUMBER);
89-
options.add(CLIENT_ID);
9080
return options;
9181
}
82+
83+
private static Configuration toFlussClientConfig(Map<String, String> tableOptions) {
84+
Configuration flussConfig = new Configuration();
85+
flussConfig.setString(
86+
ConfigOptions.BOOTSTRAP_SERVERS.key(), tableOptions.get(BOOTSTRAP_SERVERS.key()));
87+
88+
// forward all client configs
89+
tableOptions.forEach(
90+
(key, value) -> {
91+
if (key.startsWith(CLIENT_PROPERTIES_PREFIX)) {
92+
flussConfig.setString(key.substring("properties.".length()), value);
93+
}
94+
});
95+
return flussConfig;
96+
}
97+
98+
private static Map<String, String> toFlussTableOptions(Map<String, String> tableOptions) {
99+
Map<String, String> flussTableConfig = new HashMap<>();
100+
// forward all client configs
101+
tableOptions.forEach(
102+
(key, value) -> {
103+
if (key.startsWith(TABLE_PROPERTIES_PREFIX)) {
104+
flussTableConfig.put(key.substring("properties.".length()), value);
105+
}
106+
});
107+
return flussTableConfig;
108+
}
92109
}

flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/sink/FlussDataSink.java

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -31,28 +31,32 @@
3131
/** A DataSink implementation for Fluss. */
3232
public class FlussDataSink implements DataSink {
3333

34-
private final Configuration flussConfig;
34+
private final Configuration flussClientConfig;
35+
private final Map<String, String> tableProperties;
3536
private final Map<String, List<String>> bucketKeysMap;
3637
private final Map<String, Integer> bucketNumMap;
3738

3839
public FlussDataSink(
39-
Configuration flussConfig,
40+
Configuration flussClientConfig,
41+
Map<String, String> tableProperties,
4042
Map<String, List<String>> bucketKeysMap,
4143
Map<String, Integer> bucketNumMap) {
42-
this.flussConfig = flussConfig;
44+
this.flussClientConfig = flussClientConfig;
45+
this.tableProperties = tableProperties;
4346
this.bucketKeysMap = bucketKeysMap;
4447
this.bucketNumMap = bucketNumMap;
4548
}
4649

4750
@Override
4851
public EventSinkProvider getEventSinkProvider() {
4952
return FlinkSinkProvider.of(
50-
new FlinkSink<>(flussConfig, new FlussEventSerializationSchema()));
53+
new FlinkSink<>(flussClientConfig, new FlussEventSerializationSchema()));
5154
}
5255

5356
@Override
5457
public MetadataApplier getMetadataApplier() {
5558

56-
return new FlussMetaDataApplier(flussConfig, bucketKeysMap, bucketNumMap);
59+
return new FlussMetaDataApplier(
60+
flussClientConfig, tableProperties, bucketKeysMap, bucketNumMap);
5761
}
5862
}

flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/sink/FlussDataSinkOptions.java

Lines changed: 2 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,8 @@
2424
public class FlussDataSinkOptions {
2525

2626
// prefix for passing properties for fluss options.
27-
public static final String PREFIX_FLUSS_PROPERTIES = "fluss.properties.";
27+
public static final String TABLE_PROPERTIES_PREFIX = "properties.table.";
28+
public static final String CLIENT_PROPERTIES_PREFIX = "properties.client.";
2829

2930
public static final ConfigOption<String> BOOTSTRAP_SERVERS =
3031
ConfigOptions.key("bootstrap.servers")
@@ -53,10 +54,4 @@ public class FlussDataSinkOptions {
5354
"The number of buckets of each Fluss table."
5455
+ "Tables are separated by ';'. "
5556
+ "Format: database1.table1:4;database1.table2:8.");
56-
57-
public static final ConfigOption<String> CLIENT_ID =
58-
ConfigOptions.key("client.id")
59-
.stringType()
60-
.defaultValue("")
61-
.withDescription("An id string to pass to the server when making requests.");
6257
}

flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/sink/FlussMetaDataApplier.java

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -52,17 +52,20 @@
5252
/** {@link MetadataApplier} for fluss. */
5353
public class FlussMetaDataApplier implements MetadataApplier {
5454
private static final Logger LOG = LoggerFactory.getLogger(FlussMetaDataApplier.class);
55-
private final Configuration flussConfig;
55+
private final Configuration flussClientConfig;
56+
private final Map<String, String> tableProperties;
5657
private final Map<String, List<String>> bucketKeysMap;
5758
private final Map<String, Integer> bucketNumMap;
5859
private Set<SchemaChangeEventType> enabledEventTypes =
5960
new HashSet<>(Arrays.asList(CREATE_TABLE, DROP_TABLE));
6061

6162
public FlussMetaDataApplier(
62-
Configuration flussConfig,
63+
Configuration flussClientConfig,
64+
Map<String, String> tableProperties,
6365
Map<String, List<String>> bucketKeysMap,
6466
Map<String, Integer> bucketNumMap) {
65-
this.flussConfig = flussConfig;
67+
this.flussClientConfig = flussClientConfig;
68+
this.tableProperties = tableProperties;
6669
this.bucketKeysMap = bucketKeysMap;
6770
this.bucketNumMap = bucketNumMap;
6871
}
@@ -106,10 +109,10 @@ private void applyCreateTable(CreateTableEvent event) {
106109
String tableIdentifier = tablePath.getDatabaseName() + "." + tablePath.getTableName();
107110
List<String> bucketKeys = bucketKeysMap.get(tableIdentifier);
108111
Integer bucketNum = bucketNumMap.get(tableIdentifier);
109-
try (Connection connection = ConnectionFactory.createConnection(flussConfig);
112+
try (Connection connection = ConnectionFactory.createConnection(flussClientConfig);
110113
Admin admin = connection.getAdmin()) {
111114
TableDescriptor inferredFlussTable =
112-
toFlussTable(event.getSchema(), bucketKeys, bucketNum);
115+
toFlussTable(event.getSchema(), bucketKeys, bucketNum, tableProperties);
113116
admin.createDatabase(tablePath.getDatabaseName(), DatabaseDescriptor.EMPTY, true);
114117
if (!admin.tableExists(tablePath).get()) {
115118
admin.createTable(tablePath, inferredFlussTable, false).get();
@@ -126,7 +129,7 @@ private void applyCreateTable(CreateTableEvent event) {
126129
private void applyDropTable(DropTableEvent event) {
127130
TableId tableId = event.tableId();
128131
TablePath tablePath = new TablePath(tableId.getSchemaName(), tableId.getTableName());
129-
try (Connection connection = ConnectionFactory.createConnection(flussConfig);
132+
try (Connection connection = ConnectionFactory.createConnection(flussClientConfig);
130133
Admin admin = connection.getAdmin()) {
131134
admin.dropTable(tablePath, true).get();
132135
} catch (Exception e) {

flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/utils/FlinkConversions.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@
5151
import java.util.ArrayList;
5252
import java.util.Collections;
5353
import java.util.List;
54+
import java.util.Map;
5455
import java.util.stream.Collectors;
5556

5657
/** Converter from Flink's type to Fluss's type. */
@@ -60,7 +61,8 @@ public class FlinkConversions {
6061
public static TableDescriptor toFlussTable(
6162
org.apache.flink.cdc.common.schema.Schema cdcSchema,
6263
List<String> bucketKeys,
63-
@Nullable Integer bucketNum) {
64+
@Nullable Integer bucketNum,
65+
Map<String, String> tableProperties) {
6466
// now, build Fluss's table
6567
Schema.Builder schemBuilder = Schema.newBuilder();
6668
if (!CollectionUtil.isNullOrEmpty(cdcSchema.primaryKeys())) {
@@ -99,7 +101,7 @@ public static TableDescriptor toFlussTable(
99101
.partitionedBy(cdcSchema.partitionKeys())
100102
.distributedBy(bucketNum, bucketKeys)
101103
.comment(cdcSchema.comment())
102-
.properties(Collections.emptyMap())
104+
.properties(tableProperties)
103105
.build();
104106
}
105107

0 commit comments

Comments
 (0)