Skip to content

Commit 0f98f4d

Browse files
committed
add pipline support for kafka sink connector ,issue:datavane/tis#484
1 parent ff31e9d commit 0f98f4d

17 files changed

Lines changed: 400 additions & 100 deletions

File tree

.claude/settings.local.json

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,9 @@
11
{
22
"permissions": {
33
"allow": [
4-
"Bash(find /opt/misc/flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-kingbase-cdc -name *.java -type f)"
4+
"Bash(find /opt/misc/flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-kingbase-cdc -name *.java -type f)",
5+
"Bash(find /Users/mozhenghua/j2ee_solution/project/plugins/tis-datax/tis-datax-kafka-plugin -name *.java -type f)",
6+
"Bash(git -C /Users/mozhenghua/j2ee_solution/project/plugins diff HEAD -- tis-datax/tis-datax-kafka-plugin/src/main/java/com/qlangtech/tis/plugins/datax/kafka/writer/DataXKafkaWriter.java)"
57
]
68
}
79
}

tis-datax/tis-datax-kafka-plugin/src/main/java/com/qlangtech/tis/plugins/datax/kafka/writer/DataXKafkaWriter.java

Lines changed: 34 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@
1919
package com.qlangtech.tis.plugins.datax.kafka.writer;
2020

2121
import com.alibaba.citrus.turbine.Context;
22+
import com.alibaba.datax.core.job.ISourceTable;
2223
import com.google.common.collect.ImmutableMap;
24+
import com.google.common.collect.Maps;
2325
import com.qlangtech.tis.TIS;
2426
import com.qlangtech.tis.datax.IDataxContext;
2527
import com.qlangtech.tis.datax.IDataxProcessor;
@@ -33,8 +35,12 @@
3335
import com.qlangtech.tis.plugin.annotation.FormFieldType;
3436
import com.qlangtech.tis.plugin.annotation.Validator;
3537
import com.qlangtech.tis.plugin.datax.SelectedTab;
38+
import com.qlangtech.tis.plugin.datax.common.AutoCreateTable;
39+
import com.qlangtech.tis.plugin.datax.common.impl.NoneCreateTable;
3640
import com.qlangtech.tis.plugin.datax.transformer.RecordTransformerRules;
41+
import com.qlangtech.tis.plugin.ds.IInitWriterTableExecutor;
3742
import com.qlangtech.tis.plugins.datax.kafka.writer.protocol.KafkaProtocol;
43+
import com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.ICDCPipelineTableOptionsCreator;
3844
import com.qlangtech.tis.realtime.transfer.DTO;
3945
import com.qlangtech.tis.realtime.utils.NetUtils;
4046
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
@@ -51,17 +57,19 @@
5157
import org.slf4j.LoggerFactory;
5258

5359
import java.net.UnknownHostException;
60+
import java.util.List;
5461
import java.util.Map;
5562
import java.util.Objects;
5663
import java.util.Optional;
5764
import java.util.UUID;
65+
import java.util.function.Function;
5866
import java.util.stream.Collectors;
5967

6068
/**
6169
* https://nightlies.apache.org/flink/flink-docs-release-1.16/docs/connectors/datastream/kafka/ <br/>
6270
* https://github.com/airbytehq/airbyte/blob/master/airbyte-integrations/connectors/destination-kafka/src/main/resources/spec.json <br/>
6371
*/
64-
public class DataXKafkaWriter extends DataxWriter {
72+
public class DataXKafkaWriter extends DataxWriter implements IInitWriterTableExecutor, ICDCPipelineTableOptionsCreator {
6573

6674
private static final Logger LOGGER = LoggerFactory.getLogger(DataXKafkaWriter.class);
6775

@@ -138,6 +146,7 @@ public class DataXKafkaWriter extends DataxWriter {
138146
@FormField(ordinal = 22, type = FormFieldType.INT_NUMBER, validate = {Validator.require, Validator.integer}, advance = true)
139147
public Integer receiveBufferBytes;
140148

149+
141150
@Override
142151
public void startScanDependency() {
143152

@@ -186,12 +195,35 @@ public Map<String, Object> buildKafkaConfig(boolean isTest) {
186195
// .put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName())//
187196
.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName())// .StringSerializer.class.getName())
188197
.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName())// JsonSerializer.class.getName())
198+
// .put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.ByteArraySerializer")// .StringSerializer.class.getName())
199+
// .put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.ByteArraySerializer")//
189200
.build();
190201
return props.entrySet().stream()
191202
.filter(entry -> entry.getValue() != null && StringUtils.isNotBlank(entry.getValue().toString()))
192-
.collect(Collectors.toMap((e) -> e.getKey(), (e) -> String.valueOf(e.getValue())));
203+
.collect(Collectors.toMap(Map.Entry::getKey, (e) -> String.valueOf(e.getValue())));
204+
}
205+
206+
// @Override
207+
// public EntityName parseEntity(ISelectedTab tab) {
208+
// return EntityName.create(IEndTypeGetter.EndType.Kafka.getVal(), tab.getName());
209+
// }
210+
211+
@Override
212+
public AutoCreateTable getAutoCreateTableCanNotBeNull() {
213+
return new NoneCreateTable();
193214
}
194215

216+
@Override
217+
public void initWriterTable(ISourceTable sourceTable, String sinkTargetTabName, List<String> jdbcUrls) throws Exception {
218+
throw new UnsupportedOperationException();
219+
}
220+
221+
@Override
222+
public Function<SelectedTab, Map<String, String>> createTabOpts() {
223+
return (tab) -> {
224+
return Maps.newHashMap();
225+
};
226+
}
195227

196228
@TISExtension
197229
public static class DefaultDescriptor extends BaseDataxWriterDescriptor implements DataxWriter.IRewriteSuFormProperties {

tis-datax/tis-datax-kafka-plugin/src/main/java/com/qlangtech/tis/plugins/datax/kafka/writer/KafkaSelectedTab.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
import com.qlangtech.tis.util.IPluginContext;
3131
import com.qlangtech.tis.util.impl.AttrVals;
3232

33+
import java.util.Collection;
3334
import java.util.List;
3435

3536
/**
@@ -42,7 +43,7 @@ public class KafkaSelectedTab extends SelectedTab {
4243
public List<String> partitionFields;
4344

4445
public static List<Option> getPtCandidateFields() {
45-
return SelectedTab.getContextOpts((cols) -> cols.stream());
46+
return SelectedTab.getContextOpts(Collection::stream);
4647
}
4748

4849

tis-incr/tis-flink-chunjun-kafka-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/cdc/pipeline/kafka/sink/KafkaPipelineEventSinkFunc.java

Lines changed: 161 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,50 @@
11
package com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.kafka.sink;
22

3+
import com.google.common.collect.Maps;
34
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
45
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
56
import com.qlangtech.tis.datax.IDataxProcessor;
7+
import com.qlangtech.tis.datax.IDataxReader;
8+
import com.qlangtech.tis.plugin.ds.DBConfig;
9+
import com.qlangtech.tis.plugin.ds.IDataSourceFactoryGetter;
610
import com.qlangtech.tis.plugin.ds.ISelectedTab;
711
import com.qlangtech.tis.plugins.datax.kafka.writer.DataXKafkaWriter;
812
import com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.PipelineEventSinkFunc;
913
import com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.PipelineFlinkCDCSinkFactory;
14+
import com.qlangtech.tis.sql.parser.tuple.creator.EntityName;
15+
import org.apache.flink.api.common.serialization.SerializationSchema;
16+
import org.apache.flink.cdc.common.configuration.Configuration;
17+
import org.apache.flink.cdc.common.event.Event;
1018
import org.apache.flink.cdc.common.factories.Factory;
19+
import org.apache.flink.cdc.common.factories.FactoryHelper;
20+
import org.apache.flink.cdc.common.pipeline.PipelineOptions;
1121
import org.apache.flink.cdc.common.sink.DataSink;
22+
import org.apache.flink.cdc.connectors.kafka.json.ChangeLogJsonFormatFactory;
23+
import org.apache.flink.cdc.connectors.kafka.json.JsonSerializationType;
24+
import org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSink;
25+
import org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSinkFactory;
26+
import org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSinkOptions;
27+
import org.apache.flink.cdc.connectors.kafka.sink.KeyFormat;
28+
import org.apache.flink.cdc.connectors.kafka.sink.KeySerializationFactory;
29+
import org.apache.flink.cdc.connectors.kafka.sink.PartitionStrategy;
30+
import org.apache.flink.connector.base.DeliveryGuarantee;
31+
import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableMap;
32+
import org.apache.kafka.clients.producer.ProducerConfig;
1233

1334
import java.time.ZoneId;
1435
import java.util.List;
36+
import java.util.Map;
37+
import java.util.Objects;
1538
import java.util.Optional;
39+
import java.util.Properties;
40+
41+
import static org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSinkOptions.KEY_FORMAT;
42+
import static org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSinkOptions.PROPERTIES_PREFIX;
43+
import static org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSinkOptions.SINK_TABLE_ID_TO_TOPIC_MAPPING;
1644

1745
/**
1846
* <a href="https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/docs/connectors/pipeline-connectors/kafka/">...</a>
47+
*
1948
* @author 百岁 (baisui@qlangtech.com)
2049
* @date 2026/3/13
2150
* // @see tis-flink-pipeline-paimon-plugin:com.qlangtech.tis.plugins.incr.flink.pipeline.paimon.sink.PaimonPipelineEventSinkFunc
@@ -38,37 +67,149 @@ public KafkaPipelineEventSinkFunc(IDataxProcessor dataxProcessor
3867
, List<ISelectedTab> tabs //
3968
, IFlinkColCreator<FlinkCol> sourceFlinkColCreator //
4069
, int sinkTaskParallelism) {
41-
super(dataxProcessor, pipelineSinkFactory, sinkDBName, tabs, sourceFlinkColCreator, null, sinkTaskParallelism);
70+
super(dataxProcessor, pipelineSinkFactory //
71+
, sinkDBName, tabs, sourceFlinkColCreator, null, sinkTaskParallelism);
72+
}
73+
74+
@Override
75+
protected EntityName getWriterEntityName(ISelectedTab selTab) {
76+
77+
IDataxReader reader = this.dataxProcessor.getReader(null);
78+
if (reader instanceof IDataSourceFactoryGetter) {
79+
try {
80+
String[] databaseName = new String[1];
81+
DBConfig dbConfig = ((IDataSourceFactoryGetter) reader).getDataSourceFactory().getDbConfig();
82+
dbConfig.vistDbName(new DBConfig.IProcess() {
83+
@Override
84+
public boolean visit(DBConfig config, String jdbcUrl, String ip, String dbName) throws Exception {
85+
databaseName[0] = dbName;
86+
return true;
87+
}
88+
});
89+
return EntityName.create(databaseName[0], selTab.getName());
90+
} catch (Exception e) {
91+
throw new RuntimeException(e);
92+
}
93+
}
94+
95+
return super.getWriterEntityName(selTab);
4296
}
4397

4498
/**
4599
* @param context
46100
* @return
47-
* @see org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSink
101+
* @see org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSinkFactory#createDataSink(Factory.Context)
48102
*/
103+
// protected DataSink createDataSink(Factory.Context context) {
104+
// KafkaDataSinkFactory factory = new KafkaDataSinkFactory();
105+
// return factory.createDataSink(context);
106+
// }
49107
@Override
50-
protected DataSink createDataSink(Factory.Context context) {
51-
// DeliveryGuarantee deliveryGuarantee,
52-
// Properties kafkaProperties,
53-
// PartitionStrategy partitionStrategy,
54-
// ZoneId zoneId,
55-
// SerializationSchema<Event> keySerialization,
56-
// SerializationSchema<Event> valueSerialization,
57-
// String topic,
58-
// boolean addTableToHeaderEnabled,
59-
// String customHeaders,
60-
// String tableMapping
61-
62-
String topic = this.writer.topic;
63-
ZoneId zoneId = this.pipelineSinkFactory.getTimeZone();
64-
// return new KafkaDataSink(catalogOpts, tableOptions, commitUser, partitionMaps, serializer, zoneId,
65-
// schemaOperatorUid);
66-
return null;
108+
public DataSink createDataSink(Factory.Context context) {
109+
KafkaPipelineSinkFactory pipelineSinkFactory = ((KafkaPipelineSinkFactory) this.pipelineSinkFactory);
110+
KeyFormat keyFormat = context.getFactoryConfiguration().get(KEY_FORMAT);
111+
JsonSerializationType jsonSerializationType =
112+
context.getFactoryConfiguration().get(KafkaDataSinkOptions.VALUE_FORMAT);
113+
KafkaDataSinkFactory factory = new KafkaDataSinkFactory();
114+
FactoryHelper helper = FactoryHelper.createFactoryHelper(factory, context);
115+
helper.validateExcept(
116+
PROPERTIES_PREFIX, keyFormat.toString(), jsonSerializationType.toString());
117+
118+
DeliveryGuarantee deliveryGuarantee = pipelineSinkFactory.deliveryGuarantee();
119+
// context.getFactoryConfiguration().get(KafkaDataSinkOptions.DELIVERY_GUARANTEE);
120+
ZoneId zoneId = ZoneId.systemDefault();
121+
if (!Objects.equals(
122+
context.getPipelineConfiguration().get(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE),
123+
PipelineOptions.PIPELINE_LOCAL_TIME_ZONE.defaultValue())) {
124+
zoneId =
125+
ZoneId.of(
126+
context.getPipelineConfiguration()
127+
.get(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE));
128+
}
129+
SerializationSchema<Event> keySerialization =
130+
KeySerializationFactory.createSerializationSchema(
131+
helper.getFormatConfig(keyFormat.toString()), keyFormat, zoneId);
132+
org.apache.flink.configuration.Configuration formatCfg = pipelineSinkFactory.format.createFlinkCfg();
133+
134+
SerializationSchema<Event> valueSerialization =
135+
ChangeLogJsonFormatFactory.createSerializationSchema(
136+
formatCfg, // helper.getFormatConfig(jsonSerializationType.toString()),
137+
jsonSerializationType,
138+
zoneId);
139+
final Properties kafkaProperties = new Properties();
140+
Map<String, String> allOptions = context.getFactoryConfiguration().toMap();
141+
allOptions.keySet().stream()
142+
.filter(key -> key.startsWith(PROPERTIES_PREFIX))
143+
.forEach(
144+
key -> {
145+
final String value = allOptions.get(key);
146+
final String subKey = key.substring((PROPERTIES_PREFIX).length());
147+
kafkaProperties.put(subKey, value);
148+
});
149+
String topic = context.getFactoryConfiguration().get(KafkaDataSinkOptions.TOPIC);
150+
boolean addTableToHeaderEnabled =
151+
context.getFactoryConfiguration()
152+
.get(KafkaDataSinkOptions.SINK_ADD_TABLEID_TO_HEADER_ENABLED);
153+
String customHeaders =
154+
context.getFactoryConfiguration().get(KafkaDataSinkOptions.SINK_CUSTOM_HEADER);
155+
PartitionStrategy partitionStrategy = pipelineSinkFactory.getPartitionStrategy();
156+
// context.getFactoryConfiguration().get(KafkaDataSinkOptions.PARTITION_STRATEGY);
157+
String tableMapping = context.getFactoryConfiguration().get(SINK_TABLE_ID_TO_TOPIC_MAPPING);
158+
return new KafkaDataSink(
159+
deliveryGuarantee,
160+
kafkaProperties,
161+
partitionStrategy,
162+
zoneId,
163+
keySerialization,
164+
valueSerialization,
165+
topic,
166+
addTableToHeaderEnabled,
167+
customHeaders,
168+
tableMapping);
67169
}
68170

69171
@Override
70172
public Factory.Context createDataSinkContext() {
71-
return null;
173+
ImmutableMap.Builder<String, String> factoryParamBuilder = ImmutableMap.<String, String>builder();
174+
175+
// 1. topic from DataXKafkaWriter
176+
factoryParamBuilder.put(KafkaDataSinkOptions.TOPIC.key(), this.writer.topic);
177+
178+
// 2. value format from KafkaPipelineSinkFactory.format
179+
KafkaPipelineSinkFactory sinkFactory = (KafkaPipelineSinkFactory) this.pipelineSinkFactory;
180+
factoryParamBuilder.put(KafkaDataSinkOptions.VALUE_FORMAT.key(),
181+
Objects.requireNonNull(sinkFactory.format, "format can not be null").getDescriptor().getDisplayName());
182+
183+
// 3. kafka producer properties with "properties." prefix from DataXKafkaWriter
184+
Map<String, Object> kafkaConfig = this.writer.buildKafkaConfig();
185+
kafkaConfig = Maps.newHashMap(kafkaConfig);
186+
// kafkaConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,);
187+
// kafkaConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,);
188+
189+
kafkaConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG
190+
, "org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.ByteArraySerializer");// .StringSerializer.class.getName())
191+
kafkaConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG
192+
, "org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.ByteArraySerializer");//
193+
194+
for (Map.Entry<String, Object> entry : kafkaConfig.entrySet()) {
195+
// if (ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG.equals(entry.getKey())
196+
// || ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG.equals(entry.getKey())) {
197+
// continue;
198+
// }
199+
factoryParamBuilder.put(KafkaDataSinkOptions.PROPERTIES_PREFIX + entry.getKey(),
200+
String.valueOf(entry.getValue()));
201+
}
202+
203+
Configuration factoryConfiguration = Configuration.fromMap(factoryParamBuilder.build());
204+
205+
// 4. pipeline configuration with timezone
206+
ImmutableMap.Builder<String, String> pipelineParamBuilder = ImmutableMap.<String, String>builder();
207+
pipelineParamBuilder.put(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE.key(),
208+
this.pipelineSinkFactory.getTimeZone().getId());
209+
Configuration pipelineConfiguration = Configuration.fromMap(pipelineParamBuilder.build());
210+
211+
ClassLoader classLoader = KafkaPipelineEventSinkFunc.class.getClassLoader();
212+
return new FactoryHelper.DefaultContext(factoryConfiguration, pipelineConfiguration, classLoader);
72213
}
73214

74215

0 commit comments

Comments
 (0)