This is an automated email from the ASF dual-hosted git repository.
luchunliang 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 7857272d08 [INLONG-12191][Sort] Support InLong Transform SDK on the
Pulsar sink pipeline (#12192)
7857272d08 is described below
commit 7857272d0856ba83df94e42adff0410e7e7cf446
Author: ChunLiang Lu <[email protected]>
AuthorDate: Wed Aug 26 15:32:24 2026 +0800
[INLONG-12191][Sort] Support InLong Transform SDK on the Pulsar sink
pipeline (#12192)
---
.../pojo/sort/dataflow/sink/PulsarSinkConfig.java | 12 +++
.../sort/standalone/sink/cls/ClsCallback.java | 4 +
.../pulsar/DefaultEvent2PulsarRecordHandler.java | 66 +++++++++---
.../sink/pulsar/IEvent2PulsarRecordHandler.java | 17 +--
.../sink/pulsar/PulsarFederationSinkContext.java | 118 ++++++++++++++++++++-
.../standalone/sink/pulsar/PulsarIdConfig.java | 3 +
.../sink/pulsar/PulsarProducerCluster.java | 58 ++++++----
7 files changed, 230 insertions(+), 48 deletions(-)
diff --git
a/inlong-common/src/main/java/org/apache/inlong/common/pojo/sort/dataflow/sink/PulsarSinkConfig.java
b/inlong-common/src/main/java/org/apache/inlong/common/pojo/sort/dataflow/sink/PulsarSinkConfig.java
index 9142019d20..478e74fe6c 100644
---
a/inlong-common/src/main/java/org/apache/inlong/common/pojo/sort/dataflow/sink/PulsarSinkConfig.java
+++
b/inlong-common/src/main/java/org/apache/inlong/common/pojo/sort/dataflow/sink/PulsarSinkConfig.java
@@ -24,8 +24,20 @@ import lombok.EqualsAndHashCode;
@Data
public class PulsarSinkConfig extends SinkConfig {
+ public static final String MESSAGE_TYPE_CSV = "csv";
+ public static final Character CSV_DEFAULT_DELIMITER = '|';
+ public static final String MESSAGE_TYPE_KV = "kv";
+ public static final Character KV_DEFAULT_ENTRYSPLITTER = '&';
+ public static final Character KV_DEFAULT_KVSPLITTER = '=';
+ public static final String MESSAGE_TYPE_JSON = "json";
+
private String pulsarTenant;
private String namespace;
private String topic;
private Integer partitionNum;
+ private String messageType;
+ private Character delimiter;
+ private Character escapeChar;
+ private Character entrySplitter;
+ private Character kvSplitter;
}
diff --git
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsCallback.java
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsCallback.java
index cef4143876..827b063f4f 100644
---
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsCallback.java
+++
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsCallback.java
@@ -20,6 +20,7 @@ package org.apache.inlong.sort.standalone.sink.cls;
import org.apache.inlong.sort.standalone.channel.ProfileEvent;
import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
+import com.google.gson.Gson;
import com.tencentcloudapi.cls.producer.Callback;
import com.tencentcloudapi.cls.producer.Result;
import com.tencentcloudapi.cls.producer.common.Attempt;
@@ -40,6 +41,7 @@ public class ClsCallback implements Callback {
private final ClsSinkContext context;
private final ProfileEvent event;
private final String topicId;
+ private Gson gson = new Gson();
/**
* Constructor.
@@ -80,6 +82,8 @@ public class ClsCallback implements Callback {
* @param result Send result.
*/
private void onFailed(Result result) {
+ LOG.error("onFail,groupId:{},data:{},result:{}",
event.getInlongGroupId(), new String(event.getBody()),
+ gson.toJson(result));
if (isRetryable(result.getReservedAttempts())) {
tx.rollback();
tx.close();
diff --git
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/DefaultEvent2PulsarRecordHandler.java
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/DefaultEvent2PulsarRecordHandler.java
index f8016f2578..229b5c05b7 100644
---
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/DefaultEvent2PulsarRecordHandler.java
+++
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/DefaultEvent2PulsarRecordHandler.java
@@ -18,15 +18,22 @@
package org.apache.inlong.sort.standalone.sink.pulsar;
import org.apache.inlong.sdk.commons.protocol.EventConstants;
+import org.apache.inlong.sdk.transform.process.TransformProcessor;
import org.apache.inlong.sort.standalone.channel.ProfileEvent;
import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
+import com.google.gson.Gson;
import org.slf4j.Logger;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.text.SimpleDateFormat;
+import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Date;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
/**
*
@@ -40,25 +47,57 @@ public class DefaultEvent2PulsarRecordHandler implements
IEvent2PulsarRecordHand
protected final ByteArrayOutputStream outMsg = new ByteArrayOutputStream();
protected final SimpleDateFormat dateFormat = new
SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS");
protected final Date currentDate = new Date();
+ protected final Gson gson = new Gson();
/**
* parse
- *
- * @param context
- * @param event
- * @return byte array
- * @throws IOException
+ *
+ * @param context pulsar federation sink context
+ * @param event raw profile event
+ * @param idConfig id config resolved from event uid
+ * @return a list of message payload byte arrays
+ * @throws IOException on any IO error
*/
@Override
- public byte[] parse(PulsarFederationSinkContext context, ProfileEvent
event)
+ public List<byte[]> parse(PulsarFederationSinkContext context,
ProfileEvent event, PulsarIdConfig idConfig)
throws IOException {
- String uid = event.getUid();
- PulsarIdConfig idConfig = context.getIdConfig(uid);
- if (idConfig == null) {
- context.addSendResultMetric(event, context.getTaskName(), false,
System.currentTimeMillis());
- LOG.error("Can not find the id config:{}", uid);
- return null;
+ TransformProcessor<String, ?> processor =
context.getTransformProcessor(idConfig.getDataFlowId());
+ if (processor != null) {
+ return this.parseByTransform(context, event, processor);
+ } else {
+ byte[] record = this.parseByBytes(event, idConfig);
+ return Arrays.asList(record);
+ }
+ }
+
+ public List<byte[]> parseByTransform(PulsarFederationSinkContext context,
ProfileEvent event,
+ TransformProcessor<String, ?> processor) throws IOException {
+ // extParams
+ Map<String, Object> extParams = new ConcurrentHashMap<>();
+ extParams.putAll(context.getSinkContext().getParameters());
+ event.getHeaders().forEach((k, v) -> extParams.put(k, v));
+ // transform
+ List<?> results = processor.transformForBytes(event.getBody(),
extParams);
+ if (results == null) {
+ return new ArrayList<>();
+ }
+ // build
+ List<byte[]> records = new ArrayList<>(results.size());
+ for (Object result : results) {
+ byte[] msgContent;
+ if (result instanceof String) {
+ msgContent = result.toString().getBytes();
+ } else if (result instanceof byte[]) {
+ msgContent = (byte[]) result;
+ } else {
+ msgContent = gson.toJson(result).getBytes();
+ }
+ records.add(msgContent);
}
+ return records;
+ }
+
+ public byte[] parseByBytes(ProfileEvent event, PulsarIdConfig idConfig)
throws IOException {
String delimiter = idConfig.getSeparator();
byte separator = (byte) delimiter.charAt(0);
outMsg.reset();
@@ -80,8 +119,7 @@ public class DefaultEvent2PulsarRecordHandler implements
IEvent2PulsarRecordHand
break;
}
outMsg.write(event.getBody());
- byte[] msgContent = outMsg.toByteArray();
- return msgContent;
+ return outMsg.toByteArray();
}
/**
diff --git
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/IEvent2PulsarRecordHandler.java
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/IEvent2PulsarRecordHandler.java
index 694b0d757f..329d466dd9 100644
---
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/IEvent2PulsarRecordHandler.java
+++
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/IEvent2PulsarRecordHandler.java
@@ -20,6 +20,7 @@ package org.apache.inlong.sort.standalone.sink.pulsar;
import org.apache.inlong.sort.standalone.channel.ProfileEvent;
import java.io.IOException;
+import java.util.List;
/**
*
@@ -28,12 +29,14 @@ import java.io.IOException;
public interface IEvent2PulsarRecordHandler {
/**
- * parse
- *
- * @param context
- * @param event
- * @return ProducerRecord
- * @throws IOException
+ * parse the event into one or more pulsar message payloads.
+ *
+ * @param context pulsar federation sink context
+ * @param event raw profile event
+ * @param idConfig id config resolved from event uid
+ * @return a list of message payload byte arrays; empty/null
means filtered
+ * @throws IOException on any IO error
*/
- byte[] parse(PulsarFederationSinkContext context, ProfileEvent event)
throws IOException;
+ List<byte[]> parse(PulsarFederationSinkContext context, ProfileEvent
event, PulsarIdConfig idConfig)
+ throws IOException;
}
diff --git
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationSinkContext.java
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationSinkContext.java
index d17139335e..b8e88c41de 100644
---
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationSinkContext.java
+++
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationSinkContext.java
@@ -19,8 +19,18 @@ package org.apache.inlong.sort.standalone.sink.pulsar;
import org.apache.inlong.common.pojo.sort.ClusterTagConfig;
import org.apache.inlong.common.pojo.sort.TaskConfig;
+import org.apache.inlong.common.pojo.sort.dataflow.DataFlowConfig;
+import org.apache.inlong.common.pojo.sort.dataflow.sink.PulsarSinkConfig;
+import org.apache.inlong.common.pojo.sort.dataflow.sink.SinkConfig;
import org.apache.inlong.common.pojo.sort.node.PulsarNodeConfig;
import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig;
+import org.apache.inlong.sdk.transform.encode.SinkEncoder;
+import org.apache.inlong.sdk.transform.encode.SinkEncoderFactory;
+import org.apache.inlong.sdk.transform.pojo.CsvSinkInfo;
+import org.apache.inlong.sdk.transform.pojo.FieldInfo;
+import org.apache.inlong.sdk.transform.pojo.KvSinkInfo;
+import org.apache.inlong.sdk.transform.pojo.MapSinkInfo;
+import org.apache.inlong.sdk.transform.process.TransformProcessor;
import org.apache.inlong.sort.standalone.channel.ProfileEvent;
import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
import org.apache.inlong.sort.standalone.config.holder.SortClusterConfigHolder;
@@ -33,7 +43,9 @@ import
org.apache.inlong.sort.standalone.metrics.audit.AuditUtils;
import org.apache.inlong.sort.standalone.sink.SinkContext;
import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
+import com.google.common.collect.ImmutableMap;
import org.apache.commons.lang3.ClassUtils;
+import org.apache.commons.lang3.StringUtils;
import org.apache.flume.Channel;
import org.apache.flume.Context;
import org.slf4j.Logger;
@@ -45,6 +57,7 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
+@SuppressWarnings("deprecation")
public class PulsarFederationSinkContext extends SinkContext {
public static final Logger LOG =
InlongLoggerFactory.getLogger(PulsarFederationSinkContext.class);
@@ -53,6 +66,9 @@ public class PulsarFederationSinkContext extends SinkContext {
private PulsarNodeConfig pulsarNodeConfig;
private CacheClusterConfig cacheClusterConfig;
+ // Map<threadId, Map<dataFlowId,TransformProcess>>
+ protected Map<Long, Map<String, TransformProcessor<String, ?>>>
transformMap = new ConcurrentHashMap<>();
+
public PulsarFederationSinkContext(String sinkName, Context context,
Channel channel) {
super(sinkName, context, channel);
}
@@ -65,8 +81,9 @@ public class PulsarFederationSinkContext extends SinkContext {
LOG.error("newSortTaskConfig is null.");
return;
}
- if ((this.taskConfig != null &&
this.taskConfig.equals(newTaskConfig))
- && (this.sortTaskConfig != null &&
this.sortTaskConfig.equals(newSortTaskConfig))) {
+ if ((newTaskConfig == null ||
StringUtils.equals(this.taskConfigJson, gson.toJson(newTaskConfig)))
+ && (newSortTaskConfig == null
+ || StringUtils.equals(this.sortTaskConfigJson,
gson.toJson(newSortTaskConfig)))) {
LOG.info("Same sortTaskConfig, do nothing.");
return;
}
@@ -78,8 +95,7 @@ public class PulsarFederationSinkContext extends SinkContext {
}
}
- this.taskConfig = newTaskConfig;
- this.sortTaskConfig = newSortTaskConfig;
+ this.replaceConfig(newTaskConfig, newSortTaskConfig);
CacheClusterConfig clusterConfig = new CacheClusterConfig();
clusterConfig.setClusterName(this.taskName);
@@ -91,7 +107,12 @@ public class PulsarFederationSinkContext extends
SinkContext {
Map<String, PulsarIdConfig> fromTaskConfig =
fromTaskConfig(taskConfig);
Map<String, PulsarIdConfig> fromSortTaskConfig =
fromSortTaskConfig(sortTaskConfig);
SortConfigMetricReporter.reportClusterDiff(clusterId, taskName,
fromTaskConfig, fromSortTaskConfig);
- idConfigMap = unifiedConfiguration ? fromTaskConfig :
fromSortTaskConfig;
+ if (unifiedConfiguration) {
+ this.idConfigMap = fromTaskConfig;
+ this.transformMap.clear();
+ } else {
+ this.idConfigMap = fromSortTaskConfig;
+ }
} catch (Throwable e) {
LOG.error(e.getMessage(), e);
}
@@ -248,4 +269,91 @@ public class PulsarFederationSinkContext extends
SinkContext {
}
return null;
}
+
+ public TransformProcessor<String, ?> getTransformProcessor(String
dataFlowId) {
+ Long threadId = Thread.currentThread().getId();
+ Map<String, TransformProcessor<String, ?>> transformProcessors =
this.transformMap.get(threadId);
+ if (transformProcessors == null) {
+ transformProcessors = reloadTransform(taskConfig);
+ this.transformMap.put(threadId, transformProcessors);
+ }
+ return transformProcessors.get(dataFlowId);
+ }
+
+ private Map<String, TransformProcessor<String, ?>>
reloadTransform(TaskConfig taskConfig) {
+ ImmutableMap.Builder<String, TransformProcessor<String, ?>> builder =
new ImmutableMap.Builder<>();
+
+ taskConfig.getClusterTagConfigs()
+ .stream()
+ .map(ClusterTagConfig::getDataFlowConfigs)
+ .flatMap(Collection::stream)
+ .forEach(flow -> {
+ if (StringUtils.isEmpty(flow.getTransformSql())) {
+ return;
+ }
+ TransformProcessor<String, ?> transformProcessor =
createTransform(flow);
+ if (transformProcessor == null) {
+ return;
+ }
+ builder.put(flow.getDataflowId(),
+ transformProcessor);
+ });
+
+ return builder.build();
+ }
+
+ private TransformProcessor<String, ?> createTransform(DataFlowConfig
dataFlowConfig) {
+ try {
+ LOG.info("try to create transform:{}", dataFlowConfig.toString());
+ return TransformProcessor.create(
+ createTransformConfig(dataFlowConfig),
+ createSourceDecoder(dataFlowConfig.getSourceConfig()),
+ createSinkEncoder(dataFlowConfig.getSinkConfig()));
+ } catch (Exception e) {
+ LOG.error("failed to reload transform of dataflow={}, ex={}",
dataFlowConfig.getDataflowId(),
+ e.getMessage(), e);
+ return null;
+ }
+ }
+
+ private SinkEncoder<?> createSinkEncoder(SinkConfig sinkConfig) {
+ if (!(sinkConfig instanceof PulsarSinkConfig)) {
+ throw new IllegalArgumentException("sinkInfo must be an instance
of PulsarSinkConfig");
+ }
+ PulsarSinkConfig sSinkConfig = (PulsarSinkConfig) sinkConfig;
+ List<FieldInfo> fieldInfos = sSinkConfig.getFieldConfigs()
+ .stream()
+ .map(config -> new FieldInfo(config.getName(),
deriveTypeConverter(config.getFormatInfo())))
+ .collect(Collectors.toList());
+
+ if (StringUtils.equalsIgnoreCase(PulsarSinkConfig.MESSAGE_TYPE_CSV,
sSinkConfig.getMessageType())) {
+ Character delimiter = sSinkConfig.getDelimiter();
+ if (delimiter == null) {
+ delimiter = PulsarSinkConfig.CSV_DEFAULT_DELIMITER;
+ }
+ CsvSinkInfo sinkInfo = new
CsvSinkInfo(sinkConfig.getEncodingType(), delimiter,
sSinkConfig.getEscapeChar(),
+ fieldInfos);
+ return SinkEncoderFactory.createCsvEncoder(sinkInfo);
+ } else if
(StringUtils.equalsIgnoreCase(PulsarSinkConfig.MESSAGE_TYPE_KV,
sSinkConfig.getMessageType())) {
+ Character entrySplitter = sSinkConfig.getEntrySplitter();
+ if (entrySplitter == null) {
+ entrySplitter = PulsarSinkConfig.KV_DEFAULT_ENTRYSPLITTER;
+ }
+ Character kvSplitter = sSinkConfig.getKvSplitter();
+ if (kvSplitter == null) {
+ kvSplitter = PulsarSinkConfig.KV_DEFAULT_KVSPLITTER;
+ }
+ KvSinkInfo sinkInfo = new KvSinkInfo(sinkConfig.getEncodingType(),
fieldInfos);
+ sinkInfo.setEntryDelimiter(entrySplitter);
+ sinkInfo.setKvDelimiter(kvSplitter);
+ return SinkEncoderFactory.createKvEncoder(sinkInfo);
+ } else if
(StringUtils.equalsIgnoreCase(PulsarSinkConfig.MESSAGE_TYPE_JSON,
sSinkConfig.getMessageType())) {
+ MapSinkInfo sinkInfo = new
MapSinkInfo(sinkConfig.getEncodingType(), fieldInfos);
+ return SinkEncoderFactory.createMapEncoder(sinkInfo);
+ } else {
+ CsvSinkInfo sinkInfo = new
CsvSinkInfo(sinkConfig.getEncodingType(),
PulsarSinkConfig.CSV_DEFAULT_DELIMITER,
+ null, fieldInfos);
+ return SinkEncoderFactory.createCsvEncoder(sinkInfo);
+ }
+ }
}
diff --git
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarIdConfig.java
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarIdConfig.java
index 6dc48d2511..79a23bea96 100644
---
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarIdConfig.java
+++
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarIdConfig.java
@@ -55,6 +55,7 @@ public class PulsarIdConfig extends IdConfig {
private String topic;
private DataTypeEnum dataType;
private DataFlowConfig dataFlowConfig;
+ private String dataFlowId;
public PulsarIdConfig(Map<String, String> idParam) {
this.inlongGroupId = idParam.get(Constants.INLONG_GROUP_ID);
@@ -64,6 +65,7 @@ public class PulsarIdConfig extends IdConfig {
this.topic = idParam.getOrDefault(Constants.TOPIC, uid);
this.dataType = DataTypeEnum
.convert(idParam.getOrDefault(PulsarIdConfig.KEY_DATA_TYPE,
DataTypeEnum.TEXT.getType()));
+ this.dataFlowId = this.uid;
}
public static PulsarIdConfig create(DataFlowConfig dataFlowConfig) {
@@ -97,6 +99,7 @@ public class PulsarIdConfig extends IdConfig {
.dataType(dataType)
.separator(separator)
.dataFlowConfig(dataFlowConfig)
+ .dataFlowId(dataFlowConfig.getDataflowId())
.build();
}
diff --git
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerCluster.java
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerCluster.java
index 4f2f80e0ab..f4f60d4207 100644
---
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerCluster.java
+++
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerCluster.java
@@ -43,11 +43,14 @@ import org.slf4j.Logger;
import java.io.IOException;
import java.util.HashMap;
+import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
/**
*
@@ -276,35 +279,46 @@ public class PulsarProducerCluster implements
LifecycleAware {
throw new IllegalStateException();
}
- // sendAsync
- byte[] sendBytes = this.handler.parse(sinkContext, profileEvent);
+ // resolve idConfig and parse into one or more message payloads
+ String uid = profileEvent.getUid();
+ PulsarIdConfig idConfig = sinkContext.getIdConfig(uid);
+ List<byte[]> records = this.handler.parse(sinkContext, profileEvent,
idConfig);
// check
- if (sendBytes == null) {
+ if (records == null || records.isEmpty()) {
tx.commit();
profileEvent.ack();
tx.close();
return true;
}
long sendTime = System.currentTimeMillis();
- CompletableFuture<MessageId> future = producer.newMessage()
- .properties(headers)
- .value(sendBytes)
- .sendAsync();
- // callback
- future.whenCompleteAsync((msgId, ex) -> {
- if (ex != null) {
- LOG.error("Send fail:{}", ex.getMessage());
- LOG.error(ex.getMessage(), ex);
- tx.rollback();
- tx.close();
- sinkContext.addSendResultMetric(profileEvent, topic, false,
sendTime);
- } else {
- tx.commit();
- tx.close();
- sinkContext.addSendResultMetric(profileEvent, topic, true,
sendTime);
- profileEvent.ack();
- }
- });
+ // send all records, commit tx after all completed
+ final int total = records.size();
+ final AtomicInteger remaining = new AtomicInteger(total);
+ final AtomicBoolean failed = new AtomicBoolean(false);
+ for (byte[] sendBytes : records) {
+ CompletableFuture<MessageId> future = producer.newMessage()
+ .properties(headers)
+ .value(sendBytes)
+ .sendAsync();
+ future.whenCompleteAsync((msgId, ex) -> {
+ if (ex != null) {
+ LOG.error("Send fail, uid:{}, topic:{}, error:{}", uid,
topic, ex.getMessage(), ex);
+ failed.compareAndSet(false, true);
+ sinkContext.addSendResultMetric(profileEvent, topic,
false, sendTime);
+ } else {
+ sinkContext.addSendResultMetric(profileEvent, topic, true,
sendTime);
+ }
+ if (remaining.decrementAndGet() == 0) {
+ if (failed.get()) {
+ tx.rollback();
+ } else {
+ tx.commit();
+ profileEvent.ack();
+ }
+ tx.close();
+ }
+ });
+ }
return true;
}