This is an automated email from the ASF dual-hosted git repository.

dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new d2c635525 [INLONG-3667][Sort-Standalone] [Feature] Sort-Standalone add 
manageable entry to stop cache consumer before node offline or upgrade (#3714)
d2c635525 is described below

commit d2c6355258d754add233ee826ffc3900d58294d6
Author: 卢春亮 <[email protected]>
AuthorDate: Sat Apr 16 00:49:00 2022 +0800

    [INLONG-3667][Sort-Standalone] [Feature] Sort-Standalone add manageable 
entry to stop cache consumer before node offline or upgrade (#3714)
---
 .../inlong/common/metric/MetricRegister.java       |  39 ++-
 .../commons/admin/AbstractAdminEventHandler.java   |  66 +++++
 .../sdk/commons/admin/AdminEventHandler.java       |  45 +++
 .../inlong/sdk/commons/admin/AdminHttpSource.java  | 323 +++++++++++++++++++++
 .../sdk/commons/admin/AdminHttpSourceHandler.java  |  50 ++++
 .../inlong/sdk/commons/admin/AdminJsonHandler.java |  99 +++++++
 .../sdk/commons/admin/AdminServiceRegister.java    |  82 ++++++
 .../apache/inlong/sdk/commons/admin/AdminTask.java | 291 +++++++++++++++++++
 .../admin/PropertiesConfigurationProvider.java     |  63 ++++
 .../sort/impl/pulsar/InLongPulsarFetcherImpl.java  |   3 +-
 inlong-sort-standalone/bin/sort-stop.sh            |  25 +-
 inlong-sort-standalone/conf/common.properties      |   6 +
 inlong-sort-standalone/conf/es/common.properties   |   5 +
 inlong-sort-standalone/conf/hive/common.properties |   5 +
 inlong-sort-standalone/pom.xml                     |   8 +
 .../metrics/TestMetricListenerRunnable.java        |   9 +
 .../standalone/metrics/TestSortMetricItemSet.java  |   9 +
 .../apache/inlong/sort/standalone/SortCluster.java |  28 +-
 .../sort/standalone/SortStandaloneApplication.java |  11 +-
 .../admin/ConsumerServiceAdminEventHandler.java    | 152 ++++++++++
 .../standalone/admin/ConsumerServiceMBean.java     |  47 +++
 .../standalone/channel/BufferQueueChannel.java     |  12 +-
 .../standalone/channel/CacheMessageRecord.java     |  65 +++++
 .../sort/standalone/channel/ProfileEvent.java      |  46 +--
 .../sort/standalone/dispatch/DispatchProfile.java  |  12 +-
 .../sort/standalone/sink/cls/ClsCallback.java      |   1 +
 .../sort/standalone/sink/cls/ClsChannelWorker.java |   6 +
 .../sink/elasticsearch/EsCallbackListener.java     |   1 +
 .../sink/elasticsearch/EsChannelWorker.java        |   1 +
 .../inlong/sort/standalone/sink/hive/HiveSink.java |  22 +-
 .../sink/hive/WriteHdfsFileRunnable.java           |   1 +
 .../sink/kafka/KafkaFederationSinkContext.java     |  93 +++++-
 .../sink/kafka/KafkaFederationWorker.java          |  31 +-
 .../sink/kafka/KafkaProducerCluster.java           |  55 +---
 .../sink/kafka/KafkaProducerFederation.java        |   9 +-
 .../sink/pulsar/PulsarFederationSinkContext.java   |  97 ++++++-
 .../sink/pulsar/PulsarFederationWorker.java        |  43 ++-
 .../sink/pulsar/PulsarProducerCluster.java         |  46 +--
 .../sink/pulsar/PulsarProducerFederation.java      |  18 +-
 .../standalone/source/sortsdk/FetchCallback.java   |   7 +-
 .../standalone/source/sortsdk/SortSdkSource.java   |  37 ++-
 .../sink/cls/TestDefaultEvent2LogItemHandler.java  |   2 +-
 .../sink/elasticsearch/TestEsSinkContext.java      |   2 +-
 .../source/sortsdk/TestSortSdkSource.java          |  15 +-
 44 files changed, 1766 insertions(+), 222 deletions(-)

diff --git 
a/inlong-common/src/main/java/org/apache/inlong/common/metric/MetricRegister.java
 
b/inlong-common/src/main/java/org/apache/inlong/common/metric/MetricRegister.java
index 72790f2e2..4b49d0b19 100644
--- 
a/inlong-common/src/main/java/org/apache/inlong/common/metric/MetricRegister.java
+++ 
b/inlong-common/src/main/java/org/apache/inlong/common/metric/MetricRegister.java
@@ -20,9 +20,10 @@ package org.apache.inlong.common.metric;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.lang.management.ManagementFactory;
+
 import javax.management.MBeanServer;
 import javax.management.ObjectName;
-import java.lang.management.ManagementFactory;
 
 /**
  * MetricRegister
@@ -67,4 +68,40 @@ public class MetricRegister {
             LOGGER.error(ex.getMessage(), ex);
         }
     }
+
+    /**
+     * unregister MetricItem
+     */
+    public static void unregister(MetricItem obj) {
+        final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
+        StringBuilder beanName = new StringBuilder();
+        
beanName.append(JMX_DOMAIN).append(MetricItemMBean.DOMAIN_SEPARATOR).append("type=")
+                
.append(MetricUtils.getDomain(obj.getClass())).append(",").append(obj.getDimensionsKey());
+        String strBeanName = beanName.toString();
+        try {
+            ObjectName objName = new ObjectName(strBeanName);
+            mbs.unregisterMBean(objName);
+        } catch (Exception ex) {
+            LOGGER.error("exception while unregister mbean:{},error:{}", 
strBeanName, ex.getMessage());
+            LOGGER.error(ex.getMessage(), ex);
+        }
+    }
+
+    /**
+     * unregister MetricItemSet
+     */
+    public static void unregister(MetricItemSet<? extends MetricItem> obj) {
+        final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
+        StringBuilder beanName = new StringBuilder();
+        
beanName.append(JMX_DOMAIN).append(MetricItemMBean.DOMAIN_SEPARATOR).append("type=")
+                
.append(MetricUtils.getDomain(obj.getClass())).append(",name=").append(obj.getName());
+        String strBeanName = beanName.toString();
+        try {
+            ObjectName objName = new ObjectName(strBeanName);
+            mbs.unregisterMBean(objName);
+        } catch (Exception ex) {
+            LOGGER.error("exception while unregister mbean:{},error:{}", 
strBeanName, ex.getMessage());
+            LOGGER.error(ex.getMessage(), ex);
+        }
+    }
 }
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AbstractAdminEventHandler.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AbstractAdminEventHandler.java
new file mode 100644
index 000000000..405f31348
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AbstractAdminEventHandler.java
@@ -0,0 +1,66 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import java.io.IOException;
+import java.nio.charset.Charset;
+
+import javax.servlet.ServletOutputStream;
+import javax.servlet.http.HttpServletResponse;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * 
+ * AbstractAdminEventHandler
+ */
+public abstract class AbstractAdminEventHandler implements AdminEventHandler {
+
+    public static final Logger LOG = 
LoggerFactory.getLogger(AbstractAdminEventHandler.class);
+
+    /**
+     * outputResponse
+     * 
+     * @param response
+     * @param outputString
+     */
+    public void outputResponse(HttpServletResponse response, String 
outputString) {
+        ServletOutputStream outputStream = null;
+        try {
+            outputStream = response.getOutputStream();
+            
outputStream.write(outputString.getBytes(Charset.defaultCharset()));
+        } catch (Exception e) {
+            LOG.error(e.getMessage(), e);
+        } finally {
+            if (outputStream != null) {
+                try {
+                    outputStream.flush();
+                } catch (IOException e) {
+                    LOG.error(e.getMessage(), e);
+                }
+                try {
+                    outputStream.close();
+                } catch (IOException e) {
+                    LOG.error(e.getMessage(), e);
+                }
+            }
+        }
+    }
+
+}
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminEventHandler.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminEventHandler.java
new file mode 100644
index 000000000..f5d3380ad
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminEventHandler.java
@@ -0,0 +1,45 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import org.apache.flume.Event;
+import org.apache.flume.conf.Configurable;
+
+import javax.servlet.http.HttpServletResponse;
+
+/**
+ * IAdminEventHandler
+ */
+public interface AdminEventHandler extends Configurable {
+
+    String JMX_DOMAIN = "org.apache.inlong";
+    String JMX_TYPE = "type";
+    String JMX_NAME = "name";
+    char DOMAIN_SEPARATOR = ':';
+    char PROPERTY_SEPARATOR = ',';
+    char PROPERTY_EQUAL = '=';
+
+    /**
+     * process
+     * 
+     * @param cmd
+     * @param event
+     * @param response
+     */
+    void process(String cmd, Event event, HttpServletResponse response);
+}
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminHttpSource.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminHttpSource.java
new file mode 100644
index 000000000..984536fda
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminHttpSource.java
@@ -0,0 +1,323 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import java.io.IOException;
+import java.lang.management.ManagementFactory;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+import javax.servlet.http.HttpServlet;
+import javax.servlet.http.HttpServletRequest;
+import javax.servlet.http.HttpServletResponse;
+
+import org.apache.flume.ChannelException;
+import org.apache.flume.Context;
+import org.apache.flume.Event;
+import org.apache.flume.EventDrivenSource;
+import org.apache.flume.conf.Configurable;
+import org.apache.flume.instrumentation.SourceCounter;
+import org.apache.flume.source.SslContextAwareAbstractSource;
+import org.apache.flume.source.http.HTTPBadRequestException;
+import org.apache.flume.source.http.HTTPSource;
+import org.apache.flume.source.http.HTTPSourceConfigurationConstants;
+import org.apache.flume.tools.FlumeBeanConfigurator;
+import org.apache.flume.tools.HTTPServerConstraintUtil;
+import org.eclipse.jetty.http.HttpVersion;
+import org.eclipse.jetty.jmx.MBeanContainer;
+import org.eclipse.jetty.server.HttpConfiguration;
+import org.eclipse.jetty.server.HttpConnectionFactory;
+import org.eclipse.jetty.server.SecureRequestCustomizer;
+import org.eclipse.jetty.server.Server;
+import org.eclipse.jetty.server.ServerConnector;
+import org.eclipse.jetty.server.SslConnectionFactory;
+import org.eclipse.jetty.servlet.ServletContextHandler;
+import org.eclipse.jetty.servlet.ServletHolder;
+import org.eclipse.jetty.util.ssl.SslContextFactory;
+import org.eclipse.jetty.util.thread.QueuedThreadPool;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.base.Preconditions;
+import com.google.common.base.Throwables;
+
+/**
+ * AdminHttpSource
+ */
+public class AdminHttpSource extends SslContextAwareAbstractSource implements
+        EventDrivenSource, Configurable {
+    /*
+     * There are 2 ways of doing this: a. Have a static server instance and 
use connectors in each source which binds to
+     * the port defined for that source. b. Each source starts its own server 
instance, which binds to the source's
+     * port.
+     *
+     * b is more efficient than a because Jetty does not allow binding a 
servlet to a connector. So each request will
+     * need to go through each each of the handlers/servlet till the correct 
one is found.
+     *
+     */
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(HTTPSource.class);
+    private volatile Integer port;
+    private volatile Server srv;
+    private volatile String host;
+    private AdminHttpSourceHandler handler;
+    private SourceCounter sourceCounter;
+
+    private Context sourceContext;
+
+    /**
+     * configure
+     * 
+     * @param context
+     */
+    @Override
+    public void configure(Context context) {
+        configureSsl(context);
+        sourceContext = context;
+        try {
+            port = 
context.getInteger(HTTPSourceConfigurationConstants.CONFIG_PORT);
+            host = 
context.getString(HTTPSourceConfigurationConstants.CONFIG_BIND,
+                    HTTPSourceConfigurationConstants.DEFAULT_BIND);
+
+            Preconditions.checkState(host != null && !host.isEmpty(),
+                    "HTTPSource hostname specified is empty");
+            Preconditions.checkNotNull(port, "HTTPSource requires a port 
number to be"
+                    + " specified");
+
+            String handlerClassName = context.getString(
+                    HTTPSourceConfigurationConstants.CONFIG_HANDLER,
+                    HTTPSourceConfigurationConstants.DEFAULT_HANDLER).trim();
+
+            @SuppressWarnings("unchecked")
+            Class<? extends AdminHttpSourceHandler> clazz = (Class<? extends 
AdminHttpSourceHandler>) Class
+                    .forName(handlerClassName);
+            handler = clazz.getDeclaredConstructor().newInstance();
+
+            Map<String, String> subProps = context.getSubProperties(
+                    HTTPSourceConfigurationConstants.CONFIG_HANDLER_PREFIX);
+            handler.configure(new Context(subProps));
+        } catch (ClassNotFoundException ex) {
+            LOG.error("Error while configuring HTTPSource. Exception 
follows.", ex);
+            Throwables.propagate(ex);
+        } catch (ClassCastException ex) {
+            LOG.error("Deserializer is not an instance of HTTPSourceHandler."
+                    + "Deserializer must implement HTTPSourceHandler.");
+            Throwables.propagate(ex);
+        } catch (Exception ex) {
+            LOG.error("Error configuring HTTPSource!", ex);
+            Throwables.propagate(ex);
+        }
+        if (sourceCounter == null) {
+            sourceCounter = new SourceCounter(getName());
+        }
+    }
+
+    /**
+     * start
+     */
+    @Override
+    public void start() {
+        Preconditions.checkState(srv == null,
+                "Running HTTP Server found in source: " + getName()
+                        + " before I started one."
+                        + "Will not attempt to start.");
+        QueuedThreadPool threadPool = new QueuedThreadPool();
+        if (sourceContext.getSubProperties("QueuedThreadPool.").size() > 0) {
+            FlumeBeanConfigurator.setConfigurationFields(threadPool, 
sourceContext);
+        }
+        srv = new Server(threadPool);
+
+//Register with JMX for advanced monitoring
+        MBeanContainer mbContainer = new 
MBeanContainer(ManagementFactory.getPlatformMBeanServer());
+        srv.addEventListener(mbContainer);
+        srv.addBean(mbContainer);
+
+        HttpConfiguration httpConfiguration = new HttpConfiguration();
+        httpConfiguration.addCustomizer(new SecureRequestCustomizer());
+
+        FlumeBeanConfigurator.setConfigurationFields(httpConfiguration, 
sourceContext);
+        ServerConnector connector = 
getSslContextSupplier().get().map(sslContext -> {
+            SslContextFactory sslCtxFactory = new SslContextFactory();
+            sslCtxFactory.setSslContext(sslContext);
+            
sslCtxFactory.setExcludeProtocols(getExcludeProtocols().toArray(new 
String[]{}));
+            
sslCtxFactory.setIncludeProtocols(getIncludeProtocols().toArray(new 
String[]{}));
+            
sslCtxFactory.setExcludeCipherSuites(getExcludeCipherSuites().toArray(new 
String[]{}));
+            
sslCtxFactory.setIncludeCipherSuites(getIncludeCipherSuites().toArray(new 
String[]{}));
+
+            FlumeBeanConfigurator.setConfigurationFields(sslCtxFactory, 
sourceContext);
+
+            httpConfiguration.setSecurePort(port);
+            httpConfiguration.setSecureScheme("https");
+
+            return new ServerConnector(srv,
+                    new SslConnectionFactory(sslCtxFactory, 
HttpVersion.HTTP_1_1.asString()),
+                    new HttpConnectionFactory(httpConfiguration));
+        }).orElse(
+                new ServerConnector(srv, new 
HttpConnectionFactory(httpConfiguration)));
+
+        connector.setPort(port);
+        connector.setHost(host);
+        connector.setReuseAddress(true);
+
+        FlumeBeanConfigurator.setConfigurationFields(connector, sourceContext);
+
+        srv.addConnector(connector);
+
+        try {
+            ServletContextHandler context = new 
ServletContextHandler(ServletContextHandler.SESSIONS);
+            context.setContextPath("/");
+            srv.setHandler(context);
+
+            context.addServlet(new ServletHolder(new FlumeHTTPServlet()), "/");
+            
context.setSecurityHandler(HTTPServerConstraintUtil.enforceConstraints());
+            srv.start();
+        } catch (Exception ex) {
+            LOG.error("Error while starting HTTPSource. Exception follows.", 
ex);
+            Throwables.propagate(ex);
+        }
+        Preconditions.checkArgument(srv.isRunning());
+        sourceCounter.start();
+        super.start();
+    }
+
+    /**
+     * stop
+     */
+    @Override
+    public void stop() {
+        try {
+            srv.stop();
+            srv.join();
+            srv = null;
+        } catch (Exception ex) {
+            LOG.error("Error while stopping HTTPSource. Exception follows.", 
ex);
+        }
+        sourceCounter.stop();
+        LOG.info("Http source {} stopped. Metrics: {}", getName(), 
sourceCounter);
+    }
+
+    /**
+     * AdminHttpSource FlumeHTTPServlet
+     */
+    private class FlumeHTTPServlet extends HttpServlet {
+
+        private static final long serialVersionUID = 4891924863218790344L;
+
+        /**
+         * doPost
+         * 
+         * @param  request
+         * @param  response
+         * @throws IOException
+         */
+        @Override
+        public void doPost(HttpServletRequest request, HttpServletResponse 
response)
+                throws IOException {
+            List<Event> events = Collections.emptyList(); // create empty list
+            try {
+                events = handler.getEvents(request, response);
+            } catch (HTTPBadRequestException ex) {
+                LOG.warn("Received bad request from client. ", ex);
+                sourceCounter.incrementEventReadFail();
+                response.sendError(HttpServletResponse.SC_BAD_REQUEST,
+                        "Bad request from client. "
+                                + ex.getMessage());
+                return;
+            } catch (Exception ex) {
+                LOG.warn("Deserializer threw unexpected exception. ", ex);
+                sourceCounter.incrementEventReadFail();
+                
response.sendError(HttpServletResponse.SC_INTERNAL_SERVER_ERROR,
+                        "Deserializer threw unexpected exception. "
+                                + ex.getMessage());
+                return;
+            }
+            sourceCounter.incrementAppendBatchReceivedCount();
+            try {
+                if (events != null) {
+                    sourceCounter.addToEventReceivedCount(events.size());
+                    getChannelProcessor().processEventBatch(events);
+                }
+            } catch (ChannelException ex) {
+                LOG.warn("Error appending event to channel. "
+                        + "Channel might be full. Consider increasing the 
channel "
+                        + "capacity or make sure the sinks perform faster.", 
ex);
+                sourceCounter.incrementChannelWriteFail();
+                response.sendError(HttpServletResponse.SC_SERVICE_UNAVAILABLE,
+                        "Error appending event to channel. Channel might be 
full."
+                                + ex.getMessage());
+                return;
+            } catch (Exception ex) {
+                LOG.warn("Unexpected error appending event to channel. ", ex);
+                sourceCounter.incrementGenericProcessingFail();
+                
response.sendError(HttpServletResponse.SC_INTERNAL_SERVER_ERROR,
+                        "Unexpected error while appending event to channel. "
+                                + ex.getMessage());
+                return;
+            }
+            response.setCharacterEncoding(request.getCharacterEncoding());
+            response.setStatus(HttpServletResponse.SC_OK);
+            response.flushBuffer();
+            sourceCounter.incrementAppendBatchAcceptedCount();
+            if (events != null) {
+                sourceCounter.addToEventAcceptedCount(events.size());
+            }
+        }
+
+        /**
+         * doGet
+         * 
+         * @param  request
+         * @param  response
+         * @throws IOException
+         */
+        @Override
+        public void doGet(HttpServletRequest request, HttpServletResponse 
response)
+                throws IOException {
+            doPost(request, response);
+        }
+    }
+
+    /**
+     * configureSsl
+     * 
+     * @param context
+     */
+    @Override
+    protected void configureSsl(Context context) {
+        handleDeprecatedParameter(context, "ssl", "enableSSL");
+        handleDeprecatedParameter(context, "exclude-protocols", 
"excludeProtocols");
+        handleDeprecatedParameter(context, "keystore-password", 
"keystorePassword");
+
+        super.configureSsl(context);
+    }
+
+    /**
+     * handleDeprecatedParameter
+     * 
+     * @param context
+     * @param newParam
+     * @param oldParam
+     */
+    private void handleDeprecatedParameter(Context context, String newParam, 
String oldParam) {
+        if (!context.containsKey(newParam) && context.containsKey(oldParam)) {
+            context.put(newParam, context.getString(oldParam));
+        }
+    }
+
+}
\ No newline at end of file
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminHttpSourceHandler.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminHttpSourceHandler.java
new file mode 100644
index 000000000..cfdd1da24
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminHttpSourceHandler.java
@@ -0,0 +1,50 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import java.util.List;
+
+import javax.servlet.http.HttpServletRequest;
+import javax.servlet.http.HttpServletResponse;
+
+import org.apache.flume.Event;
+import org.apache.flume.conf.Configurable;
+import org.apache.flume.source.http.HTTPBadRequestException;
+
+/**
+ * 
+ * IAdminHttpSourceHandler
+ */
+public interface AdminHttpSourceHandler extends Configurable {
+
+    /**
+     * Takes an {@linkplain HttpServletRequest} and returns a list of Flume 
Events. If this request cannot be parsed
+     * into Flume events based on the format this method will throw an 
exception. This method may also throw an
+     * exception if there is some sort of other error.
+     * <p>
+     *
+     * @param  request                 The request to be parsed into Flume 
events.
+     * @param  response.
+     * @return                         List of Flume events generated from the 
request.
+     * @throws HTTPBadRequestException If the was not parsed correctly into an 
event because the request was not in the
+     *                                 expected format.
+     * @throws Exception               If there was an unexpected error.
+     */
+    List<Event> getEvents(HttpServletRequest request, HttpServletResponse 
response)
+            throws HTTPBadRequestException, Exception;
+}
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminJsonHandler.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminJsonHandler.java
new file mode 100644
index 000000000..b8ec05d14
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminJsonHandler.java
@@ -0,0 +1,99 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import org.apache.commons.lang.ClassUtils;
+import org.apache.flume.Context;
+import org.apache.flume.Event;
+import org.apache.flume.source.http.HTTPBadRequestException;
+import org.apache.flume.source.http.JSONHandler;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import javax.servlet.http.HttpServletRequest;
+import javax.servlet.http.HttpServletResponse;
+
+/**
+ * 
+ * AdminJsonHandler
+ */
+public class AdminJsonHandler implements AdminHttpSourceHandler {
+
+    public static final Logger LOG = 
LoggerFactory.getLogger(AdminJsonHandler.class);
+    public static final String KEY_CMD = "cmd";
+    private Context context;
+    private JSONHandler requestHandler;
+    private Map<String, AdminEventHandler> handlerMap = new HashMap<>();
+
+    /**
+     * configure
+     * 
+     * @param context
+     */
+    @Override
+    public void configure(Context context) {
+        this.context = context;
+        this.requestHandler = new JSONHandler();
+        this.requestHandler.configure(context);
+    }
+
+    @Override
+    public List<Event> getEvents(HttpServletRequest request, 
HttpServletResponse response)
+            throws HTTPBadRequestException, Exception {
+        List<Event> events = this.requestHandler.getEvents(request);
+        for (Event event : events) {
+            String cmd = event.getHeaders().get(KEY_CMD);
+            if (cmd == null) {
+                LOG.error("Invalid admin event,{} is null", KEY_CMD);
+                continue;
+            }
+            AdminEventHandler handler = this.handlerMap.get(cmd);
+            if (handler == null) {
+                String handlerType = context.getString(cmd + ".type");
+                if (handlerType == null) {
+                    LOG.error("Invalid admin event,{}:{},type is null", 
KEY_CMD, cmd);
+                    continue;
+                }
+                try {
+                    Class<?> handlerClass = ClassUtils.getClass(handlerType);
+                    Object handlerObject = 
handlerClass.getDeclaredConstructor().newInstance();
+                    if (handlerObject instanceof AdminEventHandler) {
+                        handler = (AdminEventHandler) handlerObject;
+                        Context subContext = new 
Context(context.getSubProperties(cmd + "."));
+                        handler.configure(subContext);
+                        this.handlerMap.put(cmd, handler);
+                    } else {
+                        LOG.error("Invalid admin event,{}:{},type:{} is not 
AdminEventHandler", KEY_CMD, cmd,
+                                handlerType);
+                        continue;
+                    }
+                } catch (Exception e) {
+                    LOG.error("Invalid admin event,{}:{},type:{} can not be 
created,error:{}", KEY_CMD, cmd,
+                            handlerType, e.getMessage(), e);
+                    continue;
+                }
+            }
+            handler.process(cmd, event, response);
+        }
+        return null;
+    }
+}
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminServiceRegister.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminServiceRegister.java
new file mode 100644
index 000000000..8018464a5
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminServiceRegister.java
@@ -0,0 +1,82 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.lang.management.ManagementFactory;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import static 
org.apache.inlong.sdk.commons.admin.AdminEventHandler.DOMAIN_SEPARATOR;
+import static org.apache.inlong.sdk.commons.admin.AdminEventHandler.JMX_DOMAIN;
+import static org.apache.inlong.sdk.commons.admin.AdminEventHandler.JMX_NAME;
+import static org.apache.inlong.sdk.commons.admin.AdminEventHandler.JMX_TYPE;
+import static 
org.apache.inlong.sdk.commons.admin.AdminEventHandler.PROPERTY_EQUAL;
+import static 
org.apache.inlong.sdk.commons.admin.AdminEventHandler.PROPERTY_SEPARATOR;
+
+/**
+ * AdminServiceRegister
+ */
+public class AdminServiceRegister {
+
+    public static final Logger LOG = 
LoggerFactory.getLogger(AdminServiceRegister.class);
+
+    /**
+     * register AdminService
+     */
+    public static void register(String type, String name, Object mbean) {
+        final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
+        String beanName = JMX_DOMAIN + DOMAIN_SEPARATOR
+                + JMX_TYPE + PROPERTY_EQUAL + type + PROPERTY_SEPARATOR
+                + JMX_NAME + PROPERTY_EQUAL + name;
+        LOG.info("start to register mbean:{}", beanName);
+        try {
+            ObjectName objName = new ObjectName(beanName);
+            mbs.registerMBean(mbean, objName);
+            LOG.info("end to register mbean:{}", beanName);
+        } catch (Exception ex) {
+            LOG.error("exception while register mbean:{},error:{}", beanName, 
ex.getMessage(), ex);
+        }
+    }
+
+    /**
+     * main
+     * 
+     * @param args
+     */
+    public static void main(String[] args) {
+        String type = "type1";
+        String name = "name1";
+        String beanName = JMX_DOMAIN + DOMAIN_SEPARATOR
+                + JMX_TYPE + PROPERTY_EQUAL + type + PROPERTY_SEPARATOR
+                + JMX_NAME + PROPERTY_EQUAL + name;
+        try {
+            ObjectName objName = new ObjectName(beanName);
+            System.out.println(objName.toString());
+            System.out.println(objName.getCanonicalKeyPropertyListString());
+            System.out.println(objName.getCanonicalName());
+            System.out.println(objName.getKeyProperty(JMX_NAME));
+        } catch (Exception ex) {
+            LOG.error("exception while register mbean:{},error:{}", beanName, 
ex.getMessage(), ex);
+        }
+
+    }
+}
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminTask.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminTask.java
new file mode 100644
index 000000000..766e0f585
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/AdminTask.java
@@ -0,0 +1,291 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import java.io.IOException;
+import java.io.StringReader;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Map.Entry;
+import java.util.Properties;
+import java.util.concurrent.locks.ReentrantLock;
+
+import org.apache.flume.Channel;
+import org.apache.flume.Context;
+import org.apache.flume.SinkRunner;
+import org.apache.flume.SourceRunner;
+import org.apache.flume.lifecycle.LifecycleState;
+import org.apache.flume.lifecycle.LifecycleSupervisor;
+import org.apache.flume.lifecycle.LifecycleSupervisor.SupervisorPolicy;
+import org.apache.flume.node.MaterializedConfiguration;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.eventbus.Subscribe;
+
+/**
+ * 
+ * AdminTask
+ */
+public class AdminTask {
+
+    public static final Logger LOG = LoggerFactory.getLogger(AdminTask.class);
+    public static final String KEY_HOST = "adminTask.host";
+    public static final String KEY_PORT = "adminTask.port";
+    public static final String KEY_HANDLER = "adminTask.handler";
+    public static final String FLUME_ROOT = "admin";
+
+    private Context context;
+    private final LifecycleSupervisor supervisor;
+    private MaterializedConfiguration materializedConfiguration;
+    private final ReentrantLock lifecycleLock = new ReentrantLock();
+
+    /**
+     * Constructor
+     * 
+     * @param context
+     */
+    public AdminTask(Context context) {
+        this.context = context;
+        this.supervisor = new LifecycleSupervisor();
+    }
+
+    /**
+     * start
+     */
+    public void start() {
+        try {
+            Map<String, String> flumeConfiguration = 
generateFlumeConfiguration();
+            if (flumeConfiguration == null) {
+                return;
+            }
+            LOG.info("Start admin task,flumeConf:{}", flumeConfiguration);
+            PropertiesConfigurationProvider configurationProvider = new 
PropertiesConfigurationProvider(
+                    FLUME_ROOT, flumeConfiguration);
+            
this.handleConfigurationEvent(configurationProvider.getConfiguration());
+        } catch (Exception e) {
+            LOG.error(e.getMessage(), e);
+        }
+    }
+
+    /**
+     * generateFlumeConfiguration
+     * 
+     * @return             Map
+     * @throws IOException
+     */
+    private Map<String, String> generateFlumeConfiguration() throws 
IOException {
+        String host = context.getString(KEY_HOST);
+        if (host == null) {
+            LOG.error("Can not start admin task, host is null.");
+            return null;
+        }
+        String port = context.getString(KEY_PORT);
+        if (port == null) {
+            LOG.error("Can not start admin task, port is null.");
+            return null;
+        }
+        String handlerType = context.getString(KEY_HANDLER);
+        if (handlerType == null) {
+            LOG.error("Can not start admin task, handlerType is null.");
+            return null;
+        }
+        String flumeString = String.format("admin.sources=r1\n"
+                + "admin.sinks=k1\n"
+                + "admin.channels=c1\n"
+                + "admin.sources.r1.type=" + AdminHttpSource.class.getName() + 
"\n"
+                + "admin.sources.r1.bind=%s\n"
+                + "admin.sources.r1.port=%s\n"
+                + "admin.sources.r1.channels=c1\n"
+                + "admin.sources.r1.handler=%s\n"
+                + "admin.sinks.k1.type=logger\n"
+                + "admin.sinks.k1.channel=c1\n"
+                + "admin.channels.c1.type=memory\n"
+                + "admin.channels.c1.capacity=1000\n"
+                + "admin.channels.c1.transactionCapacity=100", host,
+                port,
+                handlerType);
+        Properties props = new Properties();
+        props.load(new StringReader(flumeString));
+        Map<String, String> flumeMap = new HashMap<>();
+        props.forEach((key, value) -> {
+            flumeMap.put(String.valueOf(key), String.valueOf(value));
+        });
+        // 
adminTask.handler.stopService.type=org.apache.inlong.dataproxy.admin.ProxyServiceAdminEventHandler
+        // adminTask.handler.stopService.param1=xxx
+        // adminTask.handler.stopService.param2=xxx
+        // 
adminTask.handler.stopConsumer.type=org.apache.inlong.sort.standalone.admin.ConsumerServiceAdminEventHandler
+        // adminTask.handler.stopConsumer.param1=xxx
+        // adminTask.handler.stopConsumer.param2=xxx
+        // 
adminTask.handler.ackAll.type=org.apache.inlong.sort.standalone.admin.ConsumerServiceAdminEventHandler
+        // adminTask.handler.ackAll.param1=xxx
+        // adminTask.handler.ackAll.param2=xxx
+        Map<String, String> subHandlerConfig = 
context.getSubProperties(KEY_HANDLER + ".");
+        subHandlerConfig.forEach((key, value) -> {
+            flumeMap.put("admin.sources.r1.handler." + key, value);
+        });
+        return flumeMap;
+    }
+
+    /**
+     * handleConfigurationEvent
+     * 
+     * @param conf
+     */
+    @Subscribe
+    public void handleConfigurationEvent(MaterializedConfiguration conf) {
+        try {
+            lifecycleLock.lockInterruptibly();
+            stopAllComponents();
+            startAllComponents(conf);
+        } catch (InterruptedException e) {
+            LOG.info("Interrupted while trying to handle configuration event");
+            return;
+        } finally {
+            // If interrupted while trying to lock, we don't own the lock, so 
must not attempt to unlock
+            if (lifecycleLock.isHeldByCurrentThread()) {
+                lifecycleLock.unlock();
+            }
+        }
+    }
+
+    /**
+     * stop
+     */
+    public void stop() {
+        lifecycleLock.lock();
+        stopAllComponents();
+        try {
+            supervisor.stop();
+        } finally {
+            lifecycleLock.unlock();
+        }
+    }
+
+    /**
+     * stopAllComponents
+     */
+    private void stopAllComponents() {
+        if (this.materializedConfiguration != null) {
+            LOG.info("Shutting down configuration: {}", 
this.materializedConfiguration);
+            for (Entry<String, SourceRunner> entry : 
this.materializedConfiguration.getSourceRunners().entrySet()) {
+                try {
+                    LOG.info("Stopping Source " + entry.getKey());
+                    supervisor.unsupervise(entry.getValue());
+                } catch (Exception e) {
+                    LOG.error("Error while stopping {}", entry.getValue(), e);
+                }
+            }
+
+            for (Entry<String, SinkRunner> entry : 
this.materializedConfiguration.getSinkRunners().entrySet()) {
+                try {
+                    LOG.info("Stopping Sink " + entry.getKey());
+                    supervisor.unsupervise(entry.getValue());
+                } catch (Exception e) {
+                    LOG.error("Error while stopping {}", entry.getValue(), e);
+                }
+            }
+
+            for (Entry<String, Channel> entry : 
this.materializedConfiguration.getChannels().entrySet()) {
+                try {
+                    LOG.info("Stopping Channel " + entry.getKey());
+                    supervisor.unsupervise(entry.getValue());
+                } catch (Exception e) {
+                    LOG.error("Error while stopping {}", entry.getValue(), e);
+                }
+            }
+        }
+    }
+
+    /**
+     * startAllComponents
+     * 
+     * @param materializedConfiguration
+     */
+    private void startAllComponents(MaterializedConfiguration 
materializedConfiguration) {
+        LOG.info("Starting new configuration:{}", materializedConfiguration);
+
+        this.materializedConfiguration = materializedConfiguration;
+
+        for (Entry<String, Channel> entry : 
materializedConfiguration.getChannels().entrySet()) {
+            try {
+                LOG.info("Starting Channel " + entry.getKey());
+                supervisor.supervise(entry.getValue(),
+                        new SupervisorPolicy.AlwaysRestartPolicy(), 
LifecycleState.START);
+            } catch (Exception e) {
+                LOG.error("Error while starting {}", entry.getValue(), e);
+            }
+        }
+
+        /*
+         * Wait for all channels to start.
+         */
+        for (Channel ch : materializedConfiguration.getChannels().values()) {
+            while (ch.getLifecycleState() != LifecycleState.START
+                    && !supervisor.isComponentInErrorState(ch)) {
+                try {
+                    LOG.info("Waiting for channel: " + ch.getName() + " to 
start. Sleeping for 500 ms");
+                    Thread.sleep(500);
+                } catch (InterruptedException e) {
+                    LOG.error("Interrupted while waiting for channel to 
start.", e);
+                }
+            }
+        }
+
+        for (Entry<String, SinkRunner> entry : 
materializedConfiguration.getSinkRunners().entrySet()) {
+            try {
+                LOG.info("Starting Sink " + entry.getKey());
+                supervisor.supervise(entry.getValue(),
+                        new SupervisorPolicy.AlwaysRestartPolicy(), 
LifecycleState.START);
+            } catch (Exception e) {
+                LOG.error("Error while starting {}", entry.getValue(), e);
+            }
+        }
+
+        for (Entry<String, SourceRunner> entry : 
materializedConfiguration.getSourceRunners().entrySet()) {
+            try {
+                LOG.info("Starting Source " + entry.getKey());
+                supervisor.supervise(entry.getValue(),
+                        new SupervisorPolicy.AlwaysRestartPolicy(), 
LifecycleState.START);
+            } catch (Exception e) {
+                LOG.error("Error while starting {}", entry.getValue(), e);
+            }
+        }
+    }
+
+    /**
+     * main
+     * 
+     * @param args
+     */
+    public static void main(String[] args) {
+        try {
+            Context context = new Context();
+            context.put(KEY_HOST, "127.0.0.1");
+            context.put(KEY_PORT, "8080");
+            context.put(KEY_HANDLER, 
"org.apache.inlong.dataproxy.admin.AdminJsonHandler");
+            context.put(KEY_HANDLER + ".stopService.type",
+                    
"org.apache.inlong.dataproxy.admin.ProxyServiceAdminEventHandler");
+            AdminTask task = new AdminTask(context);
+            task.start();
+            Thread.sleep(10000);
+        } catch (Exception e) {
+            e.printStackTrace();
+        }
+    }
+}
\ No newline at end of file
diff --git 
a/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/PropertiesConfigurationProvider.java
 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/PropertiesConfigurationProvider.java
new file mode 100644
index 000000000..6ddfc1286
--- /dev/null
+++ 
b/inlong-sdk/sdk-common/src/main/java/org/apache/inlong/sdk/commons/admin/PropertiesConfigurationProvider.java
@@ -0,0 +1,63 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.commons.admin;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import org.apache.flume.conf.FlumeConfiguration;
+import org.apache.flume.node.AbstractConfigurationProvider;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * PropertiesConfigurationProvider
+ */
+public class PropertiesConfigurationProvider extends
+        AbstractConfigurationProvider {
+
+    public static final Logger LOG = 
LoggerFactory.getLogger(PropertiesConfigurationProvider.class);
+
+    private final Map<String, String> flumeConf;
+
+    /**
+     * PropertiesConfigurationProvider
+     * 
+     * @param agentName
+     * @param flumeConf
+     */
+    public PropertiesConfigurationProvider(String rootName, Map<String, 
String> flumeConf) {
+        super(rootName);
+        this.flumeConf = flumeConf;
+    }
+
+    /**
+     * getFlumeConfiguration
+     * 
+     * @return
+     */
+    @Override
+    public FlumeConfiguration getFlumeConfiguration() {
+        try {
+            return new FlumeConfiguration(flumeConf);
+        } catch (Exception e) {
+            LOG.error("exception catch:" + e.getMessage(), e);
+        }
+        return new FlumeConfiguration(new HashMap<String, String>());
+    }
+}
\ No newline at end of file
diff --git 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/pulsar/InLongPulsarFetcherImpl.java
 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/pulsar/InLongPulsarFetcherImpl.java
index 53841f8e6..b257d260d 100644
--- 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/pulsar/InLongPulsarFetcherImpl.java
+++ 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/pulsar/InLongPulsarFetcherImpl.java
@@ -112,7 +112,8 @@ public class InLongPulsarFetcherImpl extends 
InLongTopicFetcher {
                 consumer.acknowledgeAsync(messageId)
                         .thenAccept(consumer -> ackSucc(msgOffset))
                         .exceptionally(exception -> {
-                            logger.error("ack fail:{} {}", inLongTopic, 
msgOffset);
+                            logger.error("ack fail:{} {},error:{}", 
+                                    inLongTopic, msgOffset, 
exception.getMessage(), exception);
                             
context.getStatManager().getStatistics(context.getConfig().getSortTaskId(),
                                     
inLongTopic.getInLongCluster().getClusterId(), inLongTopic.getTopic())
                                     .addAckFailTimes(1L);
diff --git a/inlong-sort-standalone/bin/sort-stop.sh 
b/inlong-sort-standalone/bin/sort-stop.sh
index 14e122dc2..36a2d1c38 100644
--- a/inlong-sort-standalone/bin/sort-stop.sh
+++ b/inlong-sort-standalone/bin/sort-stop.sh
@@ -19,6 +19,27 @@
 # under the License.
 #
 
-# this program kills the sort
-ps -ef |grep 
"org.apache.inlong.sort.standalone.SortStandaloneApplication"|grep "java"|grep 
-v grep|awk '{print $2}'|xargs kill -9
+cd "$(dirname "$0")"/../
+export HOST_IP=`more conf/common.properties |grep "adminTask.host"|awk -F"=" 
'{print $2}'`
+export ADMIN_PORT=`more conf/common.properties |grep "adminTask.port"|awk 
-F"=" '{print $2}'`
+if [ ${HOST_IP} ] && [ ${ADMIN_PORT} ]; then
+    curl -X POST -d'[{"headers":{"cmd":"stopService"},"body":"body"}]' 
"http://${HOST_IP}:${ADMIN_PORT}";
+    echo "stop server and sleep."
+    sleep 61s
+fi
 
+#this program kills the sort
+pidInfo=$(ps -ef | grep java |grep 
org.apache.inlong.sort.standalone.SortStandaloneApplication| grep -v grep | awk 
'{print $2}')
+echo "`date` the pid info is $pidInfo">>$logFile
+
+for pid in $pidInfo;do
+       kill $pid
+done
+
+sleep 5s
+
+#force kill
+pidInfo=$(ps -ef | grep java |grep 
org.apache.inlong.sort.standalone.SortStandaloneApplication| grep -v grep | awk 
'{print $2}')
+for pid in $pidInfo;do
+       kill -9 $pid
+done
diff --git a/inlong-sort-standalone/conf/common.properties 
b/inlong-sort-standalone/conf/common.properties
index d2dcd5451..e3f3d3176 100644
--- a/inlong-sort-standalone/conf/common.properties
+++ b/inlong-sort-standalone/conf/common.properties
@@ -36,3 +36,9 @@ sortSourceConfig.QueryConsumeConfigType=file
 
#managerUrlLoaderType=org.apache.inlong.sort.standalone.config.loader.CommonPropertiesManagerUrlLoader
 
#sortClusterConfig.managerUrl=http://${manager_ip:port}/api/inlong/manager/openapi/sort/getClusterConfig
 
#sortSourceConfig.managerUrl=http://${manager_ip:port}/api/inlong/manager/openapi/sort/getSortSource
+
+adminTask.host=127.0.0.1
+adminTask.port=8088
+adminTask.handler=org.apache.inlong.sdk.commons.admin.AdminJsonHandler
+adminTask.handler.stopService.type=org.apache.inlong.sort.standalone.admin.ConsumerServiceAdminEventHandler
+
diff --git a/inlong-sort-standalone/conf/es/common.properties 
b/inlong-sort-standalone/conf/es/common.properties
index b8679449c..4495a522a 100644
--- a/inlong-sort-standalone/conf/es/common.properties
+++ b/inlong-sort-standalone/conf/es/common.properties
@@ -36,3 +36,8 @@ sortSourceConfig.QueryConsumeConfigType=file
 
#managerUrlLoaderType=org.apache.inlong.sort.standalone.config.loader.CommonPropertiesManagerUrlLoader
 
#sortClusterConfig.managerUrl=http://${manager_ip:port}/api/inlong/manager/openapi/sort/getClusterConfig
 
#sortSourceConfig.managerUrl=http://${manager_ip:port}/api/inlong/manager/openapi/sort/getSortSource
+
+adminTask.host=127.0.0.1
+adminTask.port=8088
+adminTask.handler=org.apache.inlong.sdk.commons.admin.AdminJsonHandler
+adminTask.handler.stopService.type=org.apache.inlong.sort.standalone.admin.ConsumerServiceAdminEventHandler
diff --git a/inlong-sort-standalone/conf/hive/common.properties 
b/inlong-sort-standalone/conf/hive/common.properties
index d2dcd5451..f0e43fb0e 100644
--- a/inlong-sort-standalone/conf/hive/common.properties
+++ b/inlong-sort-standalone/conf/hive/common.properties
@@ -36,3 +36,8 @@ sortSourceConfig.QueryConsumeConfigType=file
 
#managerUrlLoaderType=org.apache.inlong.sort.standalone.config.loader.CommonPropertiesManagerUrlLoader
 
#sortClusterConfig.managerUrl=http://${manager_ip:port}/api/inlong/manager/openapi/sort/getClusterConfig
 
#sortSourceConfig.managerUrl=http://${manager_ip:port}/api/inlong/manager/openapi/sort/getSortSource
+
+adminTask.host=127.0.0.1
+adminTask.port=8088
+adminTask.handler=org.apache.inlong.sdk.commons.admin.AdminJsonHandler
+adminTask.handler.stopService.type=org.apache.inlong.sort.standalone.admin.ConsumerServiceAdminEventHandler
diff --git a/inlong-sort-standalone/pom.xml b/inlong-sort-standalone/pom.xml
index ae9ba4d1d..9b1ba9386 100644
--- a/inlong-sort-standalone/pom.xml
+++ b/inlong-sort-standalone/pom.xml
@@ -175,6 +175,10 @@
             <artifactId>hive-jdbc</artifactId>
             <version>${hive.version}</version>
             <exclusions>
+                <exclusion>
+                    <groupId>org.eclipse.jetty.orbit</groupId>
+                    <artifactId>javax.servlet</artifactId>
+                </exclusion>
                 <exclusion>
                     <groupId>org.apache.hadoop</groupId>
                     <artifactId>hadoop-mapreduce-client-core</artifactId>
@@ -219,6 +223,10 @@
                     <groupId>org.openjdk.jol</groupId>
                     <artifactId>jol-core</artifactId>
                 </exclusion>
+                <exclusion>
+                    <groupId>org.eclipse.jetty.aggregate</groupId>
+                    <artifactId>jetty-all</artifactId>
+                </exclusion>
             </exclusions>
         </dependency>
         <dependency>
diff --git 
a/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestMetricListenerRunnable.java
 
b/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestMetricListenerRunnable.java
index 0ed78b580..c5e8c9d42 100644
--- 
a/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestMetricListenerRunnable.java
+++ 
b/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestMetricListenerRunnable.java
@@ -26,6 +26,7 @@ import java.util.Map;
 import org.apache.inlong.common.metric.MetricRegister;
 import org.apache.inlong.common.metric.MetricUtils;
 import org.apache.inlong.common.metric.MetricValue;
+import org.junit.AfterClass;
 import org.junit.BeforeClass;
 import org.junit.Test;
 
@@ -77,6 +78,14 @@ public class TestMetricListenerRunnable {
         itemSink.inlongStreamId = INLONG_STREAM_ID;
         dimSink = itemSink.getDimensions();
     }
+    
+    /**
+     * setdown
+     */
+    @AfterClass
+    public static void setdown() {
+        MetricRegister.unregister(itemSet);
+    }
 
     /**
      * testResult
diff --git 
a/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestSortMetricItemSet.java
 
b/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestSortMetricItemSet.java
index 7e66c1685..4849d346d 100644
--- 
a/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestSortMetricItemSet.java
+++ 
b/inlong-sort-standalone/sort-standalone-common/src/test/java/org/apache/inlong/sort/standalone/metrics/TestSortMetricItemSet.java
@@ -34,6 +34,7 @@ import org.apache.inlong.common.metric.MetricItemSetMBean;
 import org.apache.inlong.common.metric.MetricRegister;
 import org.apache.inlong.common.metric.MetricUtils;
 import org.apache.inlong.common.metric.MetricValue;
+import org.junit.AfterClass;
 import org.junit.BeforeClass;
 import org.junit.Test;
 
@@ -81,6 +82,14 @@ public class TestSortMetricItemSet {
         itemSink.inlongStreamId = INLONG_STREAM_ID;
         dimSink = itemSink.getDimensions();
     }
+    
+    /**
+     * setdown
+     */
+    @AfterClass
+    public static void setdown() {
+        MetricRegister.unregister(itemSet);
+    }
 
     /**
      * testResult
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortCluster.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortCluster.java
index 67252e0d1..4b52006c5 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortCluster.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortCluster.java
@@ -17,7 +17,14 @@
 
 package org.apache.inlong.sort.standalone;
 
-import static 
org.apache.inlong.sort.standalone.utils.Constants.RELOAD_INTERVAL;
+import org.apache.flume.Context;
+import org.apache.inlong.common.pojo.sortstandalone.SortClusterConfig;
+import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig;
+import org.apache.inlong.sdk.commons.admin.AdminTask;
+import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
+import org.apache.inlong.sort.standalone.config.holder.SortClusterConfigHolder;
+import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
+import org.slf4j.Logger;
 
 import java.util.ArrayList;
 import java.util.Date;
@@ -28,12 +35,7 @@ import java.util.Timer;
 import java.util.TimerTask;
 import java.util.concurrent.ConcurrentHashMap;
 
-import org.apache.inlong.common.pojo.sortstandalone.SortClusterConfig;
-import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig;
-import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
-import org.apache.inlong.sort.standalone.config.holder.SortClusterConfigHolder;
-import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
-import org.slf4j.Logger;
+import static 
org.apache.inlong.sort.standalone.utils.Constants.RELOAD_INTERVAL;
 
 /**
  * 
@@ -46,6 +48,7 @@ public class SortCluster {
     private Timer reloadTimer;
     private Map<String, SortTask> taskMap = new ConcurrentHashMap<>();
     private List<SortTask> deletingTasks = new ArrayList<>();
+    private AdminTask adminTask;
 
     /**
      * start
@@ -54,6 +57,9 @@ public class SortCluster {
         try {
             this.reload();
             this.setReloadTimer();
+            // start admin task
+            this.adminTask = new AdminTask(new 
Context(CommonPropertiesHolder.get()));
+            this.adminTask.start();
         } catch (Exception e) {
             LOG.error(e.getMessage(), e);
         }
@@ -65,6 +71,14 @@ public class SortCluster {
     public void close() {
         try {
             this.reloadTimer.cancel();
+            // stop sort task
+            for (Entry<String, SortTask> entry : this.taskMap.entrySet()) {
+                entry.getValue().stop();
+            }
+            // stop admin task
+            if (this.adminTask != null) {
+                this.adminTask.stop();
+            }
         } catch (Exception e) {
             LOG.error(e.getMessage(), e);
         }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortStandaloneApplication.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortStandaloneApplication.java
index f5d0b3fe4..50d6f69ce 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortStandaloneApplication.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/SortStandaloneApplication.java
@@ -41,18 +41,19 @@ public class SortStandaloneApplication {
         LOG.info("start to sort-standalone");
         try {
             SortCluster cluster = new SortCluster();
-            //
-            cluster.start();
-            // metrics
-            MetricObserver.init(CommonPropertiesHolder.get());
-            AuditUtils.initAudit();
             Runtime.getRuntime().addShutdownHook(new 
Thread("sortstandalone-shutdown-hook") {
 
                 @Override
                 public void run() {
                     AuditUtils.sendReport();
+                    cluster.close();
                 }
             });
+            //
+            cluster.start();
+            // metrics
+            MetricObserver.init(CommonPropertiesHolder.get());
+            AuditUtils.initAudit();
             Thread.sleep(5000);
         } catch (Exception e) {
             LOG.error("A fatal error occurred while running. Exception 
follows.", e);
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/admin/ConsumerServiceAdminEventHandler.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/admin/ConsumerServiceAdminEventHandler.java
new file mode 100644
index 000000000..5820cc759
--- /dev/null
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/admin/ConsumerServiceAdminEventHandler.java
@@ -0,0 +1,152 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.standalone.admin;
+
+import org.apache.commons.lang3.ClassUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.flume.Context;
+import org.apache.flume.Event;
+import org.apache.inlong.sdk.commons.admin.AbstractAdminEventHandler;
+
+import java.lang.management.ManagementFactory;
+import java.util.Set;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectInstance;
+import javax.management.ObjectName;
+import javax.servlet.http.HttpServletResponse;
+
+import static 
org.apache.inlong.sort.standalone.admin.ConsumerServiceMBean.MBEAN_TYPE;
+import static 
org.apache.inlong.sort.standalone.admin.ConsumerServiceMBean.METHOD_RECOVERCONSUMER;
+import static 
org.apache.inlong.sort.standalone.admin.ConsumerServiceMBean.METHOD_STOPCONSUMER;
+
+/**
+ * ConsumerServiceAdminEventHandler
+ */
+public class ConsumerServiceAdminEventHandler extends 
AbstractAdminEventHandler {
+
+    /**
+     * configure
+     * 
+     * @param context
+     */
+    @Override
+    public void configure(Context context) {
+    }
+
+    /**
+     * process
+     * 
+     * @param cmd
+     * @param event
+     * @param response
+     */
+    @Override
+    public void process(String cmd, Event event, HttpServletResponse response) 
{
+        LOG.info("start to process admin task:{}", cmd);
+        String sortTaskId = 
event.getHeaders().get(ConsumerServiceMBean.KEY_TASKNAME);
+        switch (cmd) {
+            case METHOD_STOPCONSUMER :
+            case METHOD_RECOVERCONSUMER :
+                if (sortTaskId == null) {
+                    break;
+                }
+                if (StringUtils.equals(sortTaskId, 
ConsumerServiceMBean.ALL_TASKNAME)) {
+                    this.processAll(cmd, event, response);
+                } else {
+                    this.processOne(cmd, sortTaskId, response);
+                }
+                break;
+            default :
+                break;
+        }
+        LOG.info("end to process admin task:{}", cmd);
+    }
+
+    /**
+     * processOne
+     * 
+     * @param cmd
+     * @param taskName
+     * @param response
+     */
+    private void processOne(String cmd, String taskName, HttpServletResponse 
response) {
+        LOG.info("start to processOne admin task:{},sort task:{}", cmd, 
taskName);
+        StringBuilder result = new StringBuilder();
+        try {
+            String beanName = JMX_DOMAIN + DOMAIN_SEPARATOR
+                    + JMX_TYPE + PROPERTY_EQUAL + MBEAN_TYPE + 
PROPERTY_SEPARATOR
+                    + JMX_NAME + PROPERTY_EQUAL + taskName;
+
+            ObjectName objName = new ObjectName(beanName);
+            final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
+            ObjectInstance mbean = mbs.getObjectInstance(objName);
+            LOG.info("getObjectInstance for type:{},name:{},result:{}", 
MBEAN_TYPE, taskName, mbean);
+            String className = mbean.getClassName();
+            Class<?> clazz = ClassUtils.getClass(className);
+            if (ClassUtils.isAssignable(clazz, ConsumerServiceMBean.class)) {
+                mbs.invoke(mbean.getObjectName(), cmd, null, null);
+                result.append(String.format("Execute command:%s success in 
bean:%s\n",
+                        cmd, mbean.getObjectName().toString()));
+            }
+            this.outputResponse(response, result.toString());
+        } catch (Exception e) {
+            LOG.error(e.getMessage(), e);
+            result.append(e.getMessage());
+            this.outputResponse(response, result.toString());
+        }
+        LOG.info("end to processOne admin task:{},sort task:{}", cmd, 
taskName);
+    }
+
+    /**
+     * processAll
+     * 
+     * @param cmd
+     * @param event
+     * @param response
+     */
+    private void processAll(String cmd, Event event, HttpServletResponse 
response) {
+        LOG.info("start to processAll admin task:{}", cmd);
+        StringBuilder result = new StringBuilder();
+        try {
+            String beanName = JMX_DOMAIN + DOMAIN_SEPARATOR
+                    + JMX_TYPE + PROPERTY_EQUAL + MBEAN_TYPE + 
PROPERTY_SEPARATOR
+                    + "*";
+            ObjectName objName = new ObjectName(beanName.toString());
+            final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
+            Set<ObjectInstance> mbeans = mbs.queryMBeans(objName, null);
+            LOG.info("queryMBeans for type:{},result:{}", MBEAN_TYPE, mbeans);
+            for (ObjectInstance mbean : mbeans) {
+                ObjectName beanObjectName = mbean.getObjectName();
+                String className = mbean.getClassName();
+                Class<?> clazz = ClassUtils.getClass(className);
+                if (ClassUtils.isAssignable(clazz, 
ConsumerServiceMBean.class)) {
+                    mbs.invoke(mbean.getObjectName(), cmd, null, null);
+                    result.append(String.format("Execute command:%s success in 
bean:%s\n",
+                            cmd, beanObjectName.toString()));
+                }
+            }
+            this.outputResponse(response, result.toString());
+        } catch (Exception e) {
+            LOG.error(e.getMessage(), e);
+            result.append(e.getMessage());
+            this.outputResponse(response, result.toString());
+        }
+        LOG.info("end to processAll admin task:{}", cmd);
+    }
+}
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/admin/ConsumerServiceMBean.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/admin/ConsumerServiceMBean.java
new file mode 100644
index 000000000..cbdd24488
--- /dev/null
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/admin/ConsumerServiceMBean.java
@@ -0,0 +1,47 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.standalone.admin;
+
+import javax.management.MXBean;
+
+/**
+ * ConsumerServiceMBean<br>
+ * Stop cache consumer before stopping Sort-Standalone process.<br>
+ * Avoid to miss the Event data in the channel when stopping Sort-Standalone 
process immediately.<br>
+ * After Sort-Standalone send all channel data to sink target and acknowledge 
all offset,<br>
+ * Sort-Standalone can be stopped.<br>
+ */
+@MXBean
+public interface ConsumerServiceMBean {
+
+    String MBEAN_TYPE = "ConsumerService";
+    String METHOD_STOPCONSUMER = "stopConsumer";
+    String KEY_TASKNAME = "sortTaskId";
+    String ALL_TASKNAME = "*";
+    String METHOD_RECOVERCONSUMER = "recoverConsumer";
+
+    /**
+     * stopConsumer
+     */
+    void stopConsumer();
+
+    /**
+     * recoverConsumer
+     */
+    void recoverConsumer();
+}
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/BufferQueueChannel.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/BufferQueueChannel.java
index b15a396e8..52f2643d3 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/BufferQueueChannel.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/BufferQueueChannel.java
@@ -17,10 +17,7 @@
 
 package org.apache.inlong.sort.standalone.channel;
 
-import java.util.Date;
-import java.util.Timer;
-import java.util.TimerTask;
-import java.util.concurrent.atomic.AtomicLong;
+import com.google.common.base.Preconditions;
 
 import org.apache.flume.ChannelException;
 import org.apache.flume.Context;
@@ -33,7 +30,10 @@ import 
org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
 import org.apache.inlong.sort.standalone.utils.SizeSemaphore;
 import org.slf4j.Logger;
 
-import com.google.common.base.Preconditions;
+import java.util.Date;
+import java.util.Timer;
+import java.util.TimerTask;
+import java.util.concurrent.atomic.AtomicLong;
 
 /**
  * 
@@ -81,7 +81,7 @@ public class BufferQueueChannel extends AbstractChannel {
             ProfileEvent profile = (ProfileEvent) event;
             transaction.doPut(profile);
         } else {
-            ProfileEvent profile = new ProfileEvent(event.getBody(), 
event.getHeaders());
+            ProfileEvent profile = new ProfileEvent(event.getBody(), 
event.getHeaders(), null);
             transaction.doPut(profile);
         }
     }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/CacheMessageRecord.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/CacheMessageRecord.java
new file mode 100644
index 000000000..63dd51750
--- /dev/null
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/CacheMessageRecord.java
@@ -0,0 +1,65 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.standalone.channel;
+
+import org.apache.inlong.sdk.sort.api.SortClient;
+import org.apache.inlong.sdk.sort.entity.MessageRecord;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * 
+ * CacheMessageRecord
+ */
+public class CacheMessageRecord {
+
+    public static final Logger LOG = 
LoggerFactory.getLogger(CacheMessageRecord.class);
+    private final SortClient client;
+    private final String msgKey;
+    private final String offset;
+    private final AtomicInteger ackCount;
+
+    /**
+     * Constructor
+     * 
+     * @param msgRecord
+     * @param client
+     */
+    public CacheMessageRecord(MessageRecord msgRecord, SortClient client) {
+        this.msgKey = msgRecord.getMsgKey();
+        this.offset = msgRecord.getOffset();
+        this.ackCount = new AtomicInteger(msgRecord.getMsgs().size());
+        this.client = client;
+    }
+
+    /**
+     * ackMessage
+     */
+    public void ackMessage() {
+        int result = this.ackCount.decrementAndGet();
+        if (result == 0 && client != null) {
+            try {
+                client.ack(msgKey, offset);
+            } catch (Exception e) {
+                LOG.error(e.getMessage(), e);
+            }
+        }
+    }
+}
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/ProfileEvent.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/ProfileEvent.java
index 153cc25db..277e02b4d 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/ProfileEvent.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/channel/ProfileEvent.java
@@ -17,13 +17,13 @@
 
 package org.apache.inlong.sort.standalone.channel;
 
-import java.util.Map;
-
 import org.apache.commons.lang.math.NumberUtils;
 import org.apache.flume.event.SimpleEvent;
 import org.apache.inlong.sort.standalone.config.pojo.InlongId;
 import org.apache.inlong.sort.standalone.utils.Constants;
 
+import java.util.Map;
+
 /**
  * 
  * ProfileEvent
@@ -36,43 +36,26 @@ public class ProfileEvent extends SimpleEvent {
 
     private final long rawLogTime;
     private final long fetchTime;
-    private long sendTime;
+    private final CacheMessageRecord cacheRecord;
 
     /**
      * Constructor
      * 
      * @param body
      * @param headers
+     * @param cacheRecord
      */
-    public ProfileEvent(byte[] body, Map<String, String> headers) {
+    public ProfileEvent(byte[] body, Map<String, String> headers, 
CacheMessageRecord cacheRecord) {
         super.setBody(body);
         super.setHeaders(headers);
+        this.cacheRecord = cacheRecord;
         this.inlongGroupId = headers.get(Constants.INLONG_GROUP_ID);
         this.inlongStreamId = headers.get(Constants.INLONG_STREAM_ID);
         this.uid = InlongId.generateUid(inlongGroupId, inlongStreamId);
         this.fetchTime = System.currentTimeMillis();
-        this.sendTime = fetchTime;
         this.rawLogTime = 
NumberUtils.toLong(headers.get(Constants.HEADER_KEY_MSG_TIME), fetchTime);
     }
 
-    /**
-     * get sendTime
-     * 
-     * @return the sendTime
-     */
-    public long getSendTime() {
-        return sendTime;
-    }
-
-    /**
-     * set sendTime
-     * 
-     * @param sendTime the sendTime to set
-     */
-    public void setSendTime(long sendTime) {
-        this.sendTime = sendTime;
-    }
-
     /**
      * get inlongGroupId
      * 
@@ -118,4 +101,21 @@ public class ProfileEvent extends SimpleEvent {
         return uid;
     }
 
+    /**
+     * get cacheRecord
+     * 
+     * @return the cacheRecord
+     */
+    public CacheMessageRecord getCacheRecord() {
+        return cacheRecord;
+    }
+
+    /**
+     * ack
+     */
+    public void ack() {
+        if (cacheRecord != null) {
+            cacheRecord.ackMessage();
+        }
+    }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/dispatch/DispatchProfile.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/dispatch/DispatchProfile.java
index a4e924661..63c25bce9 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/dispatch/DispatchProfile.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/dispatch/DispatchProfile.java
@@ -17,11 +17,11 @@
 
 package org.apache.inlong.sort.standalone.dispatch;
 
+import org.apache.inlong.sort.standalone.channel.ProfileEvent;
+
 import java.util.ArrayList;
 import java.util.List;
 
-import org.apache.inlong.sort.standalone.channel.ProfileEvent;
-
 /**
  * 
  * DispatchProfile
@@ -171,4 +171,12 @@ public class DispatchProfile {
         return dispatchTime;
     }
 
+    /**
+     * ack
+     */
+    public void ack() {
+        this.events.forEach((event) -> {
+            event.ack();
+        });
+    }
 }
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 9f900bef4..3dcc863ac 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
@@ -64,6 +64,7 @@ public class ClsCallback implements Callback {
      */
     private void onSuccess() {
         context.addSendResultMetric(event, topicId, true, 
System.currentTimeMillis());
+        event.ack();
         tx.commit();
         tx.close();
     }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsChannelWorker.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsChannelWorker.java
index 263d83fa6..7a1841cb3 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsChannelWorker.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/cls/ClsChannelWorker.java
@@ -127,6 +127,12 @@ public class ClsChannelWorker extends Thread {
     private void send(Event rowEvent, Transaction tx) throws 
ProducerException, InterruptedException {
         ProfileEvent event = (ProfileEvent) rowEvent;
         ClsIdConfig idConfig = context.getIdConfig(event.getUid());
+        if (idConfig == null) {
+            event.ack();
+            LOG.error("There is no cls id config for uid {}, discard it", 
event.getUid());
+            context.addSendResultMetric(event, context.getTaskName(), false, 
System.currentTimeMillis());
+            return;
+        }
         event.getHeaders().put(ClsSinkContext.KEY_TOPIC_ID, 
idConfig.getTopicId());
         AsyncProducerClient client = context.getClient(idConfig.getSecretId());
         List<LogItem> record = handler.parse(context, event);
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsCallbackListener.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsCallbackListener.java
index c9fbed494..2e5419492 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsCallbackListener.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsCallbackListener.java
@@ -87,6 +87,7 @@ public class EsCallbackListener implements 
BulkProcessor.Listener {
                 context.backDispatchQueue(requestItem);
             } else {
                 context.addSendResultMetric(event, context.getTaskName(), 
true, sendTime);
+                event.ack();
             }
         }
     }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsChannelWorker.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsChannelWorker.java
index a2a47121c..4be4ba6b0 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsChannelWorker.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/elasticsearch/EsChannelWorker.java
@@ -93,6 +93,7 @@ public class EsChannelWorker extends Thread {
                 context.offerDispatchQueue(indexRequest);
             } else {
                 context.addSendFailMetric();
+                profileEvent.ack();
             }
             tx.commit();
             return;
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/HiveSink.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/HiveSink.java
index 4084fe7a7..04b742470 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/HiveSink.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/HiveSink.java
@@ -17,15 +17,7 @@
 
 package org.apache.inlong.sort.standalone.sink.hive;
 
-import java.util.HashSet;
-import java.util.Map;
-import java.util.Map.Entry;
-import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.Executors;
-import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.TimeUnit;
+import com.alibaba.fastjson.JSON;
 
 import org.apache.flume.Channel;
 import org.apache.flume.Context;
@@ -40,7 +32,15 @@ import 
org.apache.inlong.sort.standalone.dispatch.DispatchProfile;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
 import org.slf4j.Logger;
 
-import com.alibaba.fastjson.JSON;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Map.Entry;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.Executors;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
 
 /**
  * 
@@ -186,6 +186,7 @@ public class HiveSink extends AbstractSink implements 
Configurable {
                 // monitor
                 LOG.error("can not find uid:{},idConfigMap:{}", uid, 
JSON.toJSONString(context.getIdConfigMap()));
                 this.context.addSendResultMetric(dispatchProfile, uid, false, 
0);
+                dispatchProfile.ack();
                 dispatchProfile = this.dispatchQueue.poll();
                 continue;
             }
@@ -203,6 +204,7 @@ public class HiveSink extends AbstractSink implements 
Configurable {
                     LOG.error(String.format("can not connect to 
hdfsPath:%s,write file:%s,error:%s",
                             context.getHdfsPath(), strIdRootPath, 
e.getMessage()), e);
                     this.context.addSendResultMetric(dispatchProfile, uid, 
false, 0);
+                    this.dispatchQueue.offer(dispatchProfile);
                     dispatchProfile = this.dispatchQueue.poll();
                     continue;
                 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/WriteHdfsFileRunnable.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/WriteHdfsFileRunnable.java
index 9848471f9..7d4e303a6 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/WriteHdfsFileRunnable.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/hive/WriteHdfsFileRunnable.java
@@ -71,6 +71,7 @@ public class WriteHdfsFileRunnable implements Runnable {
                 }
                 output.flush();
                 context.addSendResultMetric(profile, context.getTaskName(), 
true, sendTime);
+                profile.ack();
             } catch (Exception e) {
                 LOG.error(e.getMessage(), e);
                 context.addSendResultMetric(profile, context.getTaskName(), 
false, sendTime);
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationSinkContext.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationSinkContext.java
index b82d729f6..ae59cf4c9 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationSinkContext.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationSinkContext.java
@@ -17,26 +17,32 @@
 
 package org.apache.inlong.sort.standalone.sink.kafka;
 
+import org.apache.commons.lang.math.NumberUtils;
 import org.apache.flume.Channel;
 import org.apache.flume.Context;
 import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig;
+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;
 import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
 import org.apache.inlong.sort.standalone.config.pojo.InlongId;
+import org.apache.inlong.sort.standalone.metrics.SortMetricItem;
+import org.apache.inlong.sort.standalone.metrics.audit.AuditUtils;
 import org.apache.inlong.sort.standalone.sink.SinkContext;
 import org.apache.inlong.sort.standalone.utils.Constants;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
 import org.slf4j.Logger;
 
 import java.util.ArrayList;
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
 
 /** Context of kafka sink. */
 public class KafkaFederationSinkContext extends SinkContext {
-    public static final Logger LOG =
-            InlongLoggerFactory.getLogger(KafkaFederationSinkContext.class);
+
+    public static final Logger LOG = 
InlongLoggerFactory.getLogger(KafkaFederationSinkContext.class);
 
     private Context producerContext;
     private Map<String, String> idTopicMap = new ConcurrentHashMap<>();
@@ -108,8 +114,8 @@ public class KafkaFederationSinkContext extends SinkContext 
{
     /**
      * get Topic by uid
      *
-     * @param uid uid
-     * @return topic
+     * @param  uid uid
+     * @return     topic
      */
     public String getTopic(String uid) {
         String topic = this.idTopicMap.get(uid);
@@ -118,4 +124,83 @@ public class KafkaFederationSinkContext extends 
SinkContext {
         }
         return topic;
     }
+
+    /**
+     * addSendMetric
+     * 
+     * @param currentRecord
+     * @param bid
+     */
+    public void addSendMetric(ProfileEvent currentRecord, String topic) {
+        Map<String, String> dimensions = new HashMap<>();
+        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, this.getClusterId());
+        dimensions.put(SortMetricItem.KEY_TASK_NAME, this.getTaskName());
+        // metric
+        fillInlongId(currentRecord, dimensions);
+        dimensions.put(SortMetricItem.KEY_SINK_ID, this.getSinkName());
+        dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, topic);
+        long msgTime = currentRecord.getRawLogTime();
+        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
+        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
+        SortMetricItem metricItem = 
this.getMetricItemSet().findMetricItem(dimensions);
+        long count = 1;
+        long size = currentRecord.getBody().length;
+        metricItem.sendCount.addAndGet(count);
+        metricItem.sendSize.addAndGet(size);
+    }
+
+    /**
+     * addReadFailMetric
+     */
+    public void addSendFailMetric() {
+        Map<String, String> dimensions = new HashMap<>();
+        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, this.getClusterId());
+        dimensions.put(SortMetricItem.KEY_SINK_ID, this.getSinkName());
+        long msgTime = System.currentTimeMillis();
+        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
+        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
+        SortMetricItem metricItem = 
this.getMetricItemSet().findMetricItem(dimensions);
+        metricItem.readFailCount.incrementAndGet();
+    }
+
+    /**
+     * addSendResultMetric
+     * 
+     * @param currentRecord
+     * @param topic
+     * @param result
+     * @param sendTime
+     */
+    public void addSendResultMetric(ProfileEvent currentRecord, String topic, 
boolean result, long sendTime) {
+        Map<String, String> dimensions = new HashMap<>();
+        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, this.getClusterId());
+        dimensions.put(SortMetricItem.KEY_TASK_NAME, this.getTaskName());
+        // metric
+        fillInlongId(currentRecord, dimensions);
+        dimensions.put(SortMetricItem.KEY_SINK_ID, this.getSinkName());
+        dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, topic);
+        long msgTime = currentRecord.getRawLogTime();
+        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
+        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
+        SortMetricItem metricItem = 
this.getMetricItemSet().findMetricItem(dimensions);
+        long count = 1;
+        long size = currentRecord.getBody().length;
+        if (result) {
+            metricItem.sendSuccessCount.addAndGet(count);
+            metricItem.sendSuccessSize.addAndGet(size);
+            AuditUtils.add(AuditUtils.AUDIT_ID_SEND_SUCCESS, currentRecord);
+            if (sendTime > 0) {
+                long currentTime = System.currentTimeMillis();
+                long sinkDuration = currentTime - sendTime;
+                long nodeDuration = currentTime - 
NumberUtils.toLong(Constants.HEADER_KEY_SOURCE_TIME, msgTime);
+                long wholeDuration = currentTime - msgTime;
+                metricItem.sinkDuration.addAndGet(sinkDuration * count);
+                metricItem.nodeDuration.addAndGet(nodeDuration * count);
+                metricItem.wholeDuration.addAndGet(wholeDuration * count);
+            }
+        } else {
+            metricItem.sendFailCount.addAndGet(count);
+            metricItem.sendFailSize.addAndGet(size);
+        }
+    }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationWorker.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationWorker.java
index 7ce0cf30d..97f26c026 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationWorker.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaFederationWorker.java
@@ -18,12 +18,12 @@
 package org.apache.inlong.sort.standalone.sink.kafka;
 
 import com.google.common.base.Preconditions;
+
 import org.apache.flume.Channel;
 import org.apache.flume.Event;
 import org.apache.flume.Transaction;
 import org.apache.flume.lifecycle.LifecycleState;
 import org.apache.inlong.sort.standalone.channel.ProfileEvent;
-import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
 import org.apache.inlong.sort.standalone.config.pojo.InlongId;
 import org.apache.inlong.sort.standalone.metrics.SortMetricItem;
 import org.apache.inlong.sort.standalone.utils.Constants;
@@ -110,11 +110,10 @@ public class KafkaFederationWorker extends Thread {
                     LOG.error("The type of row event is not compatible with 
ProfileEvent");
                     continue;
                 }
-                ProfileEvent event = (ProfileEvent) rowEvent;
-                this.fillTopic(event);
-                SortMetricItem.fillInlongId(event, dimensions);
-                this.reportAudit(event);
-                this.producerFederation.send(event, tx);
+                ProfileEvent profileEvent = (ProfileEvent) rowEvent;
+                String topic = this.fillTopic(profileEvent);
+                this.context.addSendMetric(profileEvent, topic);
+                this.producerFederation.send(profileEvent, tx);
             } catch (Exception e) {
                 LOG.error(e.getMessage(), e);
                 if (tx != null) {
@@ -135,7 +134,7 @@ public class KafkaFederationWorker extends Thread {
      *
      * @param event
      */
-    private void fillTopic(Event event) {
+    private String fillTopic(Event event) {
         Map<String, String> headers = event.getHeaders();
         String inlongGroupId = headers.get(Constants.INLONG_GROUP_ID);
         String inlongStreamId = headers.get(Constants.INLONG_STREAM_ID);
@@ -143,7 +142,9 @@ public class KafkaFederationWorker extends Thread {
         String topic = this.context.getTopic(uid);
         if (!StringUtils.isBlank(topic)) {
             headers.put(Constants.TOPIC, topic);
+            return topic;
         }
+        return "-";
     }
 
     /** sleepOneInterval */
@@ -154,20 +155,4 @@ public class KafkaFederationWorker extends Thread {
             LOG.error(e1.getMessage(), e1);
         }
     }
-
-    /**
-     * Report event to audit.
-     *
-     * @param event Fetched event.
-     */
-    private void reportAudit(ProfileEvent event) {
-        this.dimensions.put(
-                SortMetricItem.KEY_SINK_DATA_ID, 
event.getHeaders().get(Constants.TOPIC));
-        long msgTime = event.getRawLogTime();
-        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
-        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
-        SortMetricItem metricItem = 
this.context.getMetricItemSet().findMetricItem(dimensions);
-        metricItem.sendCount.incrementAndGet();
-        metricItem.sendSize.addAndGet(event.getBody().length);
-    }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerCluster.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerCluster.java
index 12f0cd112..556183d0b 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerCluster.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerCluster.java
@@ -18,15 +18,13 @@
 package org.apache.inlong.sort.standalone.sink.kafka;
 
 import com.google.common.base.Preconditions;
+
 import org.apache.flume.Context;
-import org.apache.flume.Event;
 import org.apache.flume.Transaction;
 import org.apache.flume.lifecycle.LifecycleAware;
 import org.apache.flume.lifecycle.LifecycleState;
-import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
+import org.apache.inlong.sort.standalone.channel.ProfileEvent;
 import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
-import org.apache.inlong.sort.standalone.metrics.SortMetricItem;
-import org.apache.inlong.sort.standalone.metrics.audit.AuditUtils;
 import org.apache.inlong.sort.standalone.utils.Constants;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
 import org.apache.kafka.clients.producer.KafkaProducer;
@@ -34,11 +32,8 @@ import org.apache.kafka.clients.producer.ProducerConfig;
 import org.apache.kafka.clients.producer.ProducerRecord;
 import org.apache.kafka.common.serialization.ByteArraySerializer;
 import org.apache.kafka.common.serialization.StringSerializer;
-import org.apache.pulsar.shade.org.apache.commons.lang.math.NumberUtils;
 import org.slf4j.Logger;
 
-import java.util.HashMap;
-import java.util.Map;
 import java.util.Properties;
 
 /** wrapper of kafka producer */
@@ -46,7 +41,7 @@ public class KafkaProducerCluster implements LifecycleAware {
     public static final Logger LOG = 
InlongLoggerFactory.getLogger(KafkaProducerCluster.class);
 
     private final String workerName;
-    private final CacheClusterConfig config;
+    protected final CacheClusterConfig config;
     private final KafkaFederationSinkContext sinkContext;
     private final Context context;
 
@@ -121,23 +116,24 @@ public class KafkaProducerCluster implements 
LifecycleAware {
     /**
      * Send data
      *
-     * @param event data to send
+     * @param profileEvent data to send
      */
-    public boolean send(Event event, Transaction tx) {
-        String topic = event.getHeaders().get(Constants.TOPIC);
-        ProducerRecord<String, byte[]> record = new ProducerRecord<>(topic, 
event.getBody());
+    public boolean send(ProfileEvent profileEvent, Transaction tx) {
+        String topic = profileEvent.getHeaders().get(Constants.TOPIC);
+        ProducerRecord<String, byte[]> record = new ProducerRecord<>(topic, 
profileEvent.getBody());
         long sendTime = System.currentTimeMillis();
         try {
             producer.send(record,
                     (metadata, ex) -> {
                         if (ex == null) {
                             tx.commit();
-                            addMetric(event, topic, true, sendTime);
+                            sinkContext.addSendResultMetric(profileEvent, 
topic, true, sendTime);
+                            profileEvent.ack();
                         } else {
                             LOG.error(String.format("send failed, topic is %s, 
partition is %s",
                                             metadata.topic(), 
metadata.partition()), ex);
                             tx.rollback();
-                            addMetric(event, topic, false, 0);
+                            sinkContext.addSendResultMetric(profileEvent, 
topic, true, sendTime);
                         }
                         tx.close();
                     });
@@ -146,7 +142,7 @@ public class KafkaProducerCluster implements LifecycleAware 
{
             tx.rollback();
             tx.close();
             LOG.error(e.getMessage(), e);
-            addMetric(event, topic, false, 0);
+            sinkContext.addSendResultMetric(profileEvent, topic, true, 
sendTime);
             return false;
         }
     }
@@ -159,33 +155,4 @@ public class KafkaProducerCluster implements 
LifecycleAware {
     public String getCacheClusterName() {
         return cacheClusterName;
     }
-
-    /**
-     * Report metrics to monitor, including the count, size and duration of 
sending, sent
-     * successfully and sent failed packet.
-     *
-     * @param currentRecord event to be reported
-     * @param topic kafka topic of event sent to
-     * @param result send result, send successfully -> true, send failed -> 
false.
-     * @param sendTime the time event sent to kafka
-     */
-    private void addMetric(Event currentRecord, String topic, boolean result, 
long sendTime) {
-        Map<String, String> dimensions = new HashMap<>();
-        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, 
this.sinkContext.getClusterId());
-        // metric
-        SortMetricItem.fillInlongId(currentRecord, dimensions);
-        dimensions.put(SortMetricItem.KEY_SINK_ID, this.cacheClusterName);
-        dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, topic);
-        long msgTime =
-                
NumberUtils.toLong(currentRecord.getHeaders().get(Constants.HEADER_KEY_MSG_TIME),
 sendTime);
-        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
-        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
-        String taskName = 
currentRecord.getHeaders().get(SortMetricItem.KEY_TASK_NAME);
-        dimensions.put(SortMetricItem.KEY_TASK_NAME, taskName);
-        SortMetricItem.reportDurations(currentRecord, result, sendTime, 
dimensions, msgTime,
-                this.sinkContext.getMetricItemSet());
-        if (result) {
-            AuditUtils.add(AuditUtils.AUDIT_ID_SEND_SUCCESS, currentRecord);
-        }
-    }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerFederation.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerFederation.java
index 3c80f4e11..9ca5c3aec 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerFederation.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/kafka/KafkaProducerFederation.java
@@ -18,8 +18,9 @@
 package org.apache.inlong.sort.standalone.sink.kafka;
 
 import com.google.common.base.Preconditions;
-import org.apache.flume.Event;
+
 import org.apache.flume.Transaction;
+import org.apache.inlong.sort.standalone.channel.ProfileEvent;
 import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
 import org.slf4j.Logger;
@@ -138,11 +139,11 @@ public class KafkaProducerFederation implements Runnable {
     /**
      * send event
      *
-     * @param event event to send
+     * @param profileEvent event to send
      * @param tx transaction
      * @return send result
      */
-    public boolean send(Event event, Transaction tx) {
+    public boolean send(ProfileEvent profileEvent, Transaction tx) {
         int currentIndex = clusterIndex.getAndIncrement();
         if (currentIndex > Integer.MAX_VALUE / 2) {
             clusterIndex.set(0);
@@ -151,7 +152,7 @@ public class KafkaProducerFederation implements Runnable {
         int currentSize = currentClusterList.size();
         int realIndex = currentIndex % currentSize;
         KafkaProducerCluster clusterProducer = 
currentClusterList.get(realIndex);
-        return clusterProducer.send(event, tx);
+        return clusterProducer.send(profileEvent, tx);
     }
 
     /** Init ScheduledExecutorService with fix reload rate {@link 
#reloadInterval}. */
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 55ab5e351..f2228056d 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
@@ -17,21 +17,27 @@
 
 package org.apache.inlong.sort.standalone.sink.pulsar;
 
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
+import org.apache.commons.lang.math.NumberUtils;
 import org.apache.flume.Channel;
 import org.apache.flume.Context;
 import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig;
+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;
 import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
 import org.apache.inlong.sort.standalone.config.pojo.InlongId;
+import org.apache.inlong.sort.standalone.metrics.SortMetricItem;
+import org.apache.inlong.sort.standalone.metrics.audit.AuditUtils;
 import org.apache.inlong.sort.standalone.sink.SinkContext;
 import org.apache.inlong.sort.standalone.utils.Constants;
-import org.slf4j.Logger;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
+import org.slf4j.Logger;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
 
 /**
  * 
@@ -119,4 +125,83 @@ public class PulsarFederationSinkContext extends 
SinkContext {
     public List<CacheClusterConfig> getCacheClusters() {
         return this.clusterConfigList;
     }
+
+    /**
+     * addSendMetric
+     * 
+     * @param currentRecord
+     * @param topic
+     */
+    public void addSendMetric(ProfileEvent currentRecord, String topic) {
+        Map<String, String> dimensions = new HashMap<>();
+        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, this.getClusterId());
+        dimensions.put(SortMetricItem.KEY_TASK_NAME, this.getTaskName());
+        // metric
+        fillInlongId(currentRecord, dimensions);
+        dimensions.put(SortMetricItem.KEY_SINK_ID, this.getSinkName());
+        dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, topic);
+        long msgTime = currentRecord.getRawLogTime();
+        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
+        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
+        SortMetricItem metricItem = 
this.getMetricItemSet().findMetricItem(dimensions);
+        long count = 1;
+        long size = currentRecord.getBody().length;
+        metricItem.sendCount.addAndGet(count);
+        metricItem.sendSize.addAndGet(size);
+    }
+
+    /**
+     * addReadFailMetric
+     */
+    public void addSendFailMetric() {
+        Map<String, String> dimensions = new HashMap<>();
+        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, this.getClusterId());
+        dimensions.put(SortMetricItem.KEY_SINK_ID, this.getSinkName());
+        long msgTime = System.currentTimeMillis();
+        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
+        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
+        SortMetricItem metricItem = 
this.getMetricItemSet().findMetricItem(dimensions);
+        metricItem.readFailCount.incrementAndGet();
+    }
+
+    /**
+     * addSendResultMetric
+     * 
+     * @param currentRecord
+     * @param topic
+     * @param result
+     * @param sendTime
+     */
+    public void addSendResultMetric(ProfileEvent currentRecord, String topic, 
boolean result, long sendTime) {
+        Map<String, String> dimensions = new HashMap<>();
+        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, this.getClusterId());
+        dimensions.put(SortMetricItem.KEY_TASK_NAME, this.getTaskName());
+        // metric
+        fillInlongId(currentRecord, dimensions);
+        dimensions.put(SortMetricItem.KEY_SINK_ID, this.getSinkName());
+        dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, topic);
+        long msgTime = currentRecord.getRawLogTime();
+        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
+        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
+        SortMetricItem metricItem = 
this.getMetricItemSet().findMetricItem(dimensions);
+        long count = 1;
+        long size = currentRecord.getBody().length;
+        if (result) {
+            metricItem.sendSuccessCount.addAndGet(count);
+            metricItem.sendSuccessSize.addAndGet(size);
+            AuditUtils.add(AuditUtils.AUDIT_ID_SEND_SUCCESS, currentRecord);
+            if (sendTime > 0) {
+                long currentTime = System.currentTimeMillis();
+                long sinkDuration = currentTime - sendTime;
+                long nodeDuration = currentTime - 
NumberUtils.toLong(Constants.HEADER_KEY_SOURCE_TIME, msgTime);
+                long wholeDuration = currentTime - msgTime;
+                metricItem.sinkDuration.addAndGet(sinkDuration * count);
+                metricItem.nodeDuration.addAndGet(nodeDuration * count);
+                metricItem.wholeDuration.addAndGet(wholeDuration * count);
+            }
+        } else {
+            metricItem.sendFailCount.addAndGet(count);
+            metricItem.sendFailSize.addAndGet(size);
+        }
+    }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationWorker.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationWorker.java
index ddd5c5bcf..a42bd7cb0 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationWorker.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarFederationWorker.java
@@ -17,22 +17,19 @@
 
 package org.apache.inlong.sort.standalone.sink.pulsar;
 
-import java.util.HashMap;
-import java.util.Map;
-
 import org.apache.commons.lang.StringUtils;
 import org.apache.flume.Channel;
 import org.apache.flume.Event;
 import org.apache.flume.Transaction;
 import org.apache.flume.lifecycle.LifecycleState;
-import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
+import org.apache.inlong.sort.standalone.channel.ProfileEvent;
 import org.apache.inlong.sort.standalone.config.pojo.InlongId;
-import org.apache.inlong.sort.standalone.metrics.SortMetricItem;
 import org.apache.inlong.sort.standalone.utils.Constants;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
-import org.apache.pulsar.shade.org.apache.commons.lang.math.NumberUtils;
 import org.slf4j.Logger;
 
+import java.util.Map;
+
 /**
  * 
  * PulsarFederationWorker
@@ -46,7 +43,6 @@ public class PulsarFederationWorker extends Thread {
 
     private PulsarProducerFederation producerFederation;
     private LifecycleState status;
-    private Map<String, String> dimensions;
 
     /**
      * Constructor
@@ -61,10 +57,6 @@ public class PulsarFederationWorker extends Thread {
         this.context = context;
         this.producerFederation = new PulsarProducerFederation(workerName, 
this.context);
         this.status = LifecycleState.IDLE;
-        this.dimensions = new HashMap<>();
-        this.dimensions.put(SortMetricItem.KEY_CLUSTER_ID, 
this.context.getClusterId());
-        this.dimensions.put(SortMetricItem.KEY_TASK_NAME, 
this.context.getTaskName());
-        this.dimensions.put(SortMetricItem.KEY_SINK_ID, 
this.context.getSinkName());
     }
 
     /**
@@ -104,28 +96,27 @@ public class PulsarFederationWorker extends Thread {
                     sleepOneInterval();
                     continue;
                 }
+                if (!(event instanceof ProfileEvent)) {
+                    tx.commit();
+                    this.context.addSendFailMetric();
+                    Thread.sleep(context.getProcessInterval());
+                    return;
+                }
+                // to profileEvent
+                ProfileEvent profileEvent = (ProfileEvent) event;
                 // fill topic
-                this.fillTopic(event);
+                String topic = this.fillTopic(profileEvent);
                 // metric
-                SortMetricItem.fillInlongId(event, dimensions);
-                this.dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, 
event.getHeaders().get(Constants.TOPIC));
-                long msgTime = 
NumberUtils.toLong(event.getHeaders().get(Constants.HEADER_KEY_MSG_TIME),
-                        System.currentTimeMillis());
-                long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
-                dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
-                SortMetricItem metricItem = 
this.context.getMetricItemSet().findMetricItem(dimensions);
-                metricItem.sendCount.incrementAndGet();
-                metricItem.sendSize.addAndGet(event.getBody().length);
+                this.context.addSendMetric(profileEvent, topic);
                 // send
-                this.producerFederation.send(event, tx);
+                this.producerFederation.send(profileEvent, tx);
             } catch (Throwable t) {
                 LOG.error("Process event failed!" + this.getName(), t);
                 try {
                     tx.rollback();
                     tx.close();
                     // metric
-                    SortMetricItem metricItem = 
this.context.getMetricItemSet().findMetricItem(dimensions);
-                    metricItem.readFailCount.incrementAndGet();
+                    context.addSendFailMetric();
                     sleepOneInterval();
                 } catch (Throwable e) {
                     LOG.error("Channel take transaction rollback exception:" + 
getName(), e);
@@ -139,7 +130,7 @@ public class PulsarFederationWorker extends Thread {
      * 
      * @param currentRecord
      */
-    private void fillTopic(Event currentRecord) {
+    private String fillTopic(Event currentRecord) {
         Map<String, String> headers = currentRecord.getHeaders();
         String inlongGroupId = headers.get(Constants.INLONG_GROUP_ID);
         String inlongStreamId = headers.get(Constants.INLONG_STREAM_ID);
@@ -147,7 +138,9 @@ public class PulsarFederationWorker extends Thread {
         String topic = this.context.getTopic(uid);
         if (!StringUtils.isBlank(topic)) {
             headers.put(Constants.TOPIC, topic);
+            return topic;
         }
+        return "-";
     }
 
     /**
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 1772a210f..ec101de3c 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
@@ -18,14 +18,11 @@
 package org.apache.inlong.sort.standalone.sink.pulsar;
 
 import org.apache.flume.Context;
-import org.apache.flume.Event;
 import org.apache.flume.Transaction;
 import org.apache.flume.lifecycle.LifecycleAware;
 import org.apache.flume.lifecycle.LifecycleState;
-import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
+import org.apache.inlong.sort.standalone.channel.ProfileEvent;
 import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
-import org.apache.inlong.sort.standalone.metrics.SortMetricItem;
-import org.apache.inlong.sort.standalone.metrics.audit.AuditUtils;
 import org.apache.inlong.sort.standalone.utils.Constants;
 import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
 import org.apache.pulsar.client.api.AuthenticationFactory;
@@ -37,10 +34,8 @@ import org.apache.pulsar.client.api.Producer;
 import org.apache.pulsar.client.api.ProducerBuilder;
 import org.apache.pulsar.client.api.PulsarClient;
 import org.apache.pulsar.client.api.PulsarClientException;
-import org.apache.pulsar.shade.org.apache.commons.lang.math.NumberUtils;
 import org.slf4j.Logger;
 
-import java.util.HashMap;
 import java.util.Map;
 import java.util.Map.Entry;
 import java.util.concurrent.CompletableFuture;
@@ -194,12 +189,12 @@ public class PulsarProducerCluster implements 
LifecycleAware {
     /**
      * send
      * 
-     * @param event
+     * @param profileEvent
      * @param tx
      */
-    public boolean send(Event event, Transaction tx) {
+    public boolean send(ProfileEvent profileEvent, Transaction tx) {
         // send
-        Map<String, String> headers = event.getHeaders();
+        Map<String, String> headers = profileEvent.getHeaders();
         String topic = headers.get(Constants.TOPIC);
         // get producer
         Producer<byte[]> producer = this.producerMap.get(topic);
@@ -224,7 +219,7 @@ public class PulsarProducerCluster implements 
LifecycleAware {
         if (producer == null) {
             tx.rollback();
             tx.close();
-            this.addMetric(event, topic, false, 0);
+            sinkContext.addSendResultMetric(profileEvent, topic, false, 
System.currentTimeMillis());
             return false;
         }
         String messageKey = headers.get(Constants.HEADER_KEY_MESSAGE_KEY);
@@ -234,7 +229,7 @@ public class PulsarProducerCluster implements 
LifecycleAware {
         // sendAsync
         long sendTime = System.currentTimeMillis();
         CompletableFuture<MessageId> future = 
producer.newMessage().key(messageKey).properties(headers)
-                .value(event.getBody()).sendAsync();
+                .value(profileEvent.getBody()).sendAsync();
         // callback
         future.whenCompleteAsync((msgId, ex) -> {
             if (ex != null) {
@@ -242,40 +237,17 @@ public class PulsarProducerCluster implements 
LifecycleAware {
                 LOG.error(ex.getMessage(), ex);
                 tx.rollback();
                 tx.close();
-                this.addMetric(event, topic, false, 0);
+                sinkContext.addSendResultMetric(profileEvent, topic, false, 
sendTime);
             } else {
                 tx.commit();
                 tx.close();
-                this.addMetric(event, topic, true, sendTime);
+                sinkContext.addSendResultMetric(profileEvent, topic, true, 
sendTime);
+                profileEvent.ack();
             }
         });
         return true;
     }
 
-    /**
-     * addMetric
-     * 
-     * @param currentRecord
-     * @param topic
-     * @param result
-     */
-    private void addMetric(Event currentRecord, String topic, boolean result, 
long sendTime) {
-        Map<String, String> dimensions = new HashMap<>();
-        dimensions.put(SortMetricItem.KEY_CLUSTER_ID, 
this.sinkContext.getClusterId());
-        // metric
-        SortMetricItem.fillInlongId(currentRecord, dimensions);
-        dimensions.put(SortMetricItem.KEY_SINK_ID, this.cacheClusterName);
-        dimensions.put(SortMetricItem.KEY_SINK_DATA_ID, topic);
-        long msgTime = 
NumberUtils.toLong(currentRecord.getHeaders().get(Constants.HEADER_KEY_MSG_TIME),
 sendTime);
-        long auditFormatTime = msgTime - msgTime % 
CommonPropertiesHolder.getAuditFormatInterval();
-        dimensions.put(SortMetricItem.KEY_MESSAGE_TIME, 
String.valueOf(auditFormatTime));
-        SortMetricItem.reportDurations(currentRecord, result, sendTime,
-                dimensions, msgTime, this.sinkContext.getMetricItemSet());
-        if (result) {
-            AuditUtils.add(AuditUtils.AUDIT_ID_SEND_SUCCESS, currentRecord);
-        }
-    }
-
     /**
      * get cacheClusterName
      * 
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerFederation.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerFederation.java
index c914fc888..19344efd5 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerFederation.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/sink/pulsar/PulsarProducerFederation.java
@@ -17,6 +17,12 @@
 
 package org.apache.inlong.sort.standalone.sink.pulsar;
 
+import org.apache.flume.Transaction;
+import org.apache.inlong.sort.standalone.channel.ProfileEvent;
+import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
+import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
+import org.slf4j.Logger;
+
 import java.util.ArrayList;
 import java.util.Date;
 import java.util.HashSet;
@@ -26,12 +32,6 @@ import java.util.Timer;
 import java.util.TimerTask;
 import java.util.concurrent.atomic.AtomicInteger;
 
-import org.apache.flume.Event;
-import org.apache.flume.Transaction;
-import org.apache.inlong.sort.standalone.config.pojo.CacheClusterConfig;
-import org.slf4j.Logger;
-import org.apache.inlong.sort.standalone.utils.InlongLoggerFactory;
-
 /**
  * 
  * PulsarProducerFederation
@@ -148,10 +148,10 @@ public class PulsarProducerFederation {
     /**
      * send
      * 
-     * @param event
+     * @param profileEvent
      * @param tx
      */
-    public boolean send(Event event, Transaction tx) {
+    public boolean send(ProfileEvent profileEvent, Transaction tx) {
         int currentIndex = clusterIndex.getAndIncrement();
         if (currentIndex > Integer.MAX_VALUE / 2) {
             clusterIndex.set(0);
@@ -160,6 +160,6 @@ public class PulsarProducerFederation {
         int currentSize = currentClusterList.size();
         int realIndex = currentIndex % currentSize;
         PulsarProducerCluster clusterProducer = 
currentClusterList.get(realIndex);
-        return clusterProducer.send(event, tx);
+        return clusterProducer.send(profileEvent, tx);
     }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/FetchCallback.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/FetchCallback.java
index edb4ddaac..2fd854efa 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/FetchCallback.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/FetchCallback.java
@@ -26,6 +26,7 @@ import org.apache.inlong.sdk.sort.api.ReadCallback;
 import org.apache.inlong.sdk.sort.api.SortClient;
 import org.apache.inlong.sdk.sort.entity.InLongMessage;
 import org.apache.inlong.sdk.sort.entity.MessageRecord;
+import org.apache.inlong.sort.standalone.channel.CacheMessageRecord;
 import org.apache.inlong.sort.standalone.channel.ProfileEvent;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -102,18 +103,20 @@ public class FetchCallback implements ReadCallback {
     public void onFinished(final MessageRecord messageRecord) {
         try {
             Preconditions.checkState(messageRecord != null, "Fetched msg is 
null.");
+            CacheMessageRecord cacheRecord = new 
CacheMessageRecord(messageRecord, client);
             for (InLongMessage inLongMessage : messageRecord.getMsgs()) {
                 //TODO fix here
                 final SubscribeFetchResult result = 
SubscribeFetchResult.Factory
                         .create(sortId, messageRecord.getMsgKey(), 
messageRecord.getOffset(),
                                 inLongMessage.getParams(), 
messageRecord.getRecTime(),
                                 inLongMessage.getBody());
-                final ProfileEvent profileEvent = new 
ProfileEvent(result.getBody(), result.getHeaders());
+                final ProfileEvent profileEvent = new 
ProfileEvent(result.getBody(), result.getHeaders(), 
+                        cacheRecord);
                 channelProcessor.processEvent(profileEvent);
                 context.reportToMetric(profileEvent, sortId, "-", 
SortSdkSourceContext.FetchResult.SUCCESS);
             }
 
-            client.ack(messageRecord.getMsgKey(), messageRecord.getOffset());
+//            client.ack(messageRecord.getMsgKey(), messageRecord.getOffset());
         } catch (NullPointerException npe) {
             LOG.error("Got a null pointer exception for sortId " + sortId, 
npe);
             context.reportToMetric(null, sortId, "-", 
SortSdkSourceContext.FetchResult.FAILURE);
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/SortSdkSource.java
 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/SortSdkSource.java
index d9f0cee2b..42b8b34ff 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/SortSdkSource.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/main/java/org/apache/inlong/sort/standalone/source/sortsdk/SortSdkSource.java
@@ -17,23 +17,19 @@
 
 package org.apache.inlong.sort.standalone.source.sortsdk;
 
-import java.net.InetAddress;
-import java.net.UnknownHostException;
-import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.TimeUnit;
-
 import org.apache.commons.lang.ClassUtils;
 import org.apache.flume.Context;
 import org.apache.flume.EventDrivenSource;
 import org.apache.flume.conf.Configurable;
 import org.apache.flume.source.AbstractSource;
+import org.apache.inlong.sdk.commons.admin.AdminServiceRegister;
 import org.apache.inlong.sdk.sort.api.QueryConsumeConfig;
 import org.apache.inlong.sdk.sort.api.SortClient;
 import org.apache.inlong.sdk.sort.api.SortClientConfig;
 import org.apache.inlong.sdk.sort.api.SortClientFactory;
 import org.apache.inlong.sdk.sort.impl.ManagerReportHandlerImpl;
 import org.apache.inlong.sdk.sort.impl.MetricReporterImpl;
+import org.apache.inlong.sort.standalone.admin.ConsumerServiceMBean;
 import org.apache.inlong.sort.standalone.config.holder.CommonPropertiesHolder;
 import org.apache.inlong.sort.standalone.config.holder.ManagerUrlHandler;
 import org.apache.inlong.sort.standalone.config.holder.SortClusterConfigHolder;
@@ -44,6 +40,12 @@ import 
org.apache.inlong.sort.standalone.utils.FlumeConfigGenerator;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.net.InetAddress;
+import java.net.UnknownHostException;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+
 /**
  * Default Source implementation of InLong.
  *
@@ -62,7 +64,8 @@ import org.slf4j.LoggerFactory;
  * <b>WITHOUT</b> any arguments, and parameters will be configured by {@link 
Configurable#configure(Context)}.
  * </p>
  */
-public final class SortSdkSource extends AbstractSource implements 
Configurable, Runnable, EventDrivenSource {
+public final class SortSdkSource extends AbstractSource
+        implements Configurable, Runnable, EventDrivenSource, 
ConsumerServiceMBean {
 
     // Log of {@link SortSdkSource}.
     private static final Logger LOG = 
LoggerFactory.getLogger(SortSdkSource.class);
@@ -107,6 +110,7 @@ public final class SortSdkSource extends AbstractSource 
implements Configurable,
         pool.shutdownNow();
         LOG.info("Close sort client {}.", taskName);
         if (sortClient != null) {
+            sortClient.getConfig().setStopConsume(true);
             sortClient.close();
         }
     }
@@ -134,7 +138,8 @@ public final class SortSdkSource extends AbstractSource 
implements Configurable,
         this.sortClusterName = 
SortClusterConfigHolder.getClusterConfig().getClusterName();
         this.reloadInterval = this.context.getReloadInterval();
         this.initReloadExecutor();
-
+        // register
+        AdminServiceRegister.register(ConsumerServiceMBean.MBEAN_TYPE, 
taskName, this);
     }
 
     /**
@@ -212,4 +217,20 @@ public final class SortSdkSource extends AbstractSource 
implements Configurable,
         }
         return null;
     }
+
+    /**
+     * stopConsumer
+     */
+    @Override
+    public void stopConsumer() {
+        sortClient.getConfig().setStopConsume(true);
+    }
+
+    /**
+     * recoverConsumer
+     */
+    @Override
+    public void recoverConsumer() {
+        sortClient.getConfig().setStopConsume(false);
+    }
 }
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/cls/TestDefaultEvent2LogItemHandler.java
 
b/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/cls/TestDefaultEvent2LogItemHandler.java
index 211c8f0d1..9e8a69792 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/cls/TestDefaultEvent2LogItemHandler.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/cls/TestDefaultEvent2LogItemHandler.java
@@ -85,7 +85,7 @@ public class TestDefaultEvent2LogItemHandler {
         headers.put(Constants.INLONG_GROUP_ID, "testGroup");
         headers.put(Constants.INLONG_STREAM_ID, "testStream");
         headers.put(Constants.HEADER_KEY_MSG_TIME, "1234456");
-        return new ProfileEvent(body, headers);
+        return new ProfileEvent(body, headers, null);
     }
 
 }
\ No newline at end of file
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/elasticsearch/TestEsSinkContext.java
 
b/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/elasticsearch/TestEsSinkContext.java
index 3982bd65d..732f78484 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/elasticsearch/TestEsSinkContext.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/sink/elasticsearch/TestEsSinkContext.java
@@ -86,7 +86,7 @@ public class TestEsSinkContext {
         headers.put(Constants.HEADER_KEY_MSG_TIME, 
String.valueOf(System.currentTimeMillis()));
         headers.put(Constants.HEADER_KEY_SOURCE_IP, "127.0.0.1");
         byte[] body = content.getBytes(Charset.defaultCharset());
-        return new ProfileEvent(body, headers);
+        return new ProfileEvent(body, headers, null);
     }
 
     /**
diff --git 
a/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/source/sortsdk/TestSortSdkSource.java
 
b/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/source/sortsdk/TestSortSdkSource.java
index 5bac874e0..89ba0429a 100644
--- 
a/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/source/sortsdk/TestSortSdkSource.java
+++ 
b/inlong-sort-standalone/sort-standalone-source/src/test/java/org/apache/inlong/sort/standalone/source/sortsdk/TestSortSdkSource.java
@@ -21,6 +21,7 @@ import org.apache.flume.Context;
 import org.apache.inlong.common.metric.MetricRegister;
 import org.apache.inlong.common.pojo.sortstandalone.SortClusterConfig;
 import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig;
+import org.apache.inlong.sdk.commons.admin.AdminServiceRegister;
 import org.apache.inlong.sort.standalone.config.holder.SortClusterConfigHolder;
 import org.junit.Assert;
 import org.junit.Before;
@@ -28,6 +29,7 @@ import org.junit.Test;
 import org.junit.runner.RunWith;
 import org.mockito.Mockito;
 import org.powermock.api.mockito.PowerMockito;
+import org.powermock.core.classloader.annotations.PowerMockIgnore;
 import org.powermock.core.classloader.annotations.PrepareForTest;
 import org.powermock.modules.junit4.PowerMockRunner;
 import org.slf4j.Logger;
@@ -36,8 +38,13 @@ import org.slf4j.LoggerFactory;
 import java.util.ArrayList;
 import java.util.List;
 
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+
 @RunWith(PowerMockRunner.class)
-@PrepareForTest({SortClusterConfigHolder.class, LoggerFactory.class, 
Logger.class, MetricRegister.class})
+@PowerMockIgnore("javax.management.*")
+@PrepareForTest({SortClusterConfigHolder.class, LoggerFactory.class, 
Logger.class, MetricRegister.class,
+        AdminServiceRegister.class})
 public class TestSortSdkSource {
 
     private Context mockContext;
@@ -57,6 +64,12 @@ public class TestSortSdkSource {
 
     @Test
     public void testRun() {
+        PowerMockito.mockStatic(AdminServiceRegister.class);
+        try {
+            PowerMockito.doNothing().when(AdminServiceRegister.class, 
"register", anyString(), anyString(), any());
+        } catch (Exception e) {
+            e.printStackTrace();
+        }
         SortSdkSource testSource = new SortSdkSource();
         testSource.configure(mockContext);
         testSource.run();

Reply via email to