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;
     }
 

Reply via email to