RockteMQ-AI commented on code in PR #384:
URL: https://github.com/apache/rocketmq-connect/pull/384#discussion_r3839551606


##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/ApacheHttpClientImpl.java:
##########
@@ -0,0 +1,263 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import com.google.common.net.MediaType;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.http.HttpEntity;
+import org.apache.http.HttpHost;
+import org.apache.http.client.config.RequestConfig;
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpDelete;
+import org.apache.http.client.methods.HttpGet;
+import org.apache.http.client.methods.HttpHead;
+import org.apache.http.client.methods.HttpOptions;
+import org.apache.http.client.methods.HttpPatch;
+import org.apache.http.client.methods.HttpPost;
+import org.apache.http.client.methods.HttpPut;
+import org.apache.http.client.methods.HttpRequestBase;
+import org.apache.http.client.methods.HttpTrace;
+import org.apache.http.client.protocol.HttpClientContext;
+import org.apache.http.config.Registry;
+import org.apache.http.config.RegistryBuilder;
+import org.apache.http.conn.DnsResolver;
+import org.apache.http.conn.socket.ConnectionSocketFactory;
+import org.apache.http.conn.socket.PlainConnectionSocketFactory;
+import org.apache.http.conn.ssl.SSLConnectionSocketFactory;
+import org.apache.http.conn.ssl.TrustStrategy;
+import org.apache.http.entity.StringEntity;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
+import org.apache.http.protocol.HTTP;
+import org.apache.http.protocol.HttpContext;
+import org.apache.http.ssl.SSLContextBuilder;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.slf4j.MDC;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.SSLSession;
+import java.io.IOException;
+import java.io.UnsupportedEncodingException;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.Proxy;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import java.security.cert.X509Certificate;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import static 
org.apache.rocketmq.connect.http.sink.constant.HttpConstant.LOG_SIFT_TAG;
+
+
+public class ApacheHttpClientImpl implements AbstractHttpClient {
+    private static final Logger log = 
LoggerFactory.getLogger(ApacheHttpClientImpl.class);
+
+    private static ExecutorService executorServicePool = new 
ThreadPoolExecutor(200, 2000, 600, TimeUnit.SECONDS,
+            new LinkedBlockingDeque<Runnable>(1000), new 
DefaultThreadFactory("ApacheHttpClientRequestThread"));
+    private CloseableHttpClient httpClient = null;
+
+    private SocksProxyConfig socksProxyConfig;
+    private static final String SOCKS_ADDRESS_KEY = "socks.address";
+
+    @Override
+    public void init(ClientConfig config) {
+        try {
+            SSLContextBuilder sslContextBuilder = new 
SSLContextBuilder().loadTrustMaterial(null, new TrustStrategy() {
+                @Override

Review Comment:
   SSL/TLS verification is completely disabled: TrustStrategy.isTrusted always 
returns true (line 77) and NoopHostnameVerifier.verify always returns true 
(line 253). This makes all HTTPS connections vulnerable to man-in-the-middle 
attacks. This should at minimum be opt-in via a configuration flag, not the 
unconditional default.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());
+                httpRequest.setHeaderMap(headerMap);
+                httpRequest.setMethod(clientConfig.getMethod());
+                httpRequest.setTimeout(clientConfig.getTimeout());
+                
httpRequest.setUrl(JsonUtils.queryStringAndPathValue(clientConfig.getUrlPattern(),
 clientConfig.getQueryStringParameters(), 
connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE)));
+                httpClient.execute(httpRequest, httpCallback);
+            }
+            boolean consumeSucceed = Boolean.FALSE;
+            try {
+                consumeSucceed = 
countDownLatch.await(DEFAULT_CONSUMER_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+            } catch (Throwable e) {
+                log.error("count down latch failed.", e);
+            }
+            if (!consumeSucceed) {
+                throw new RetriableException("Request Timeout");
+            }
+            if (httpCallback.isFailed()) {
+                throw new RetriableException(httpCallback.getMsg());
+            }
         } catch (Exception e) {
             log.error("HttpSinkTask | put | error => ", e);
+            throw new RuntimeException(e);
         }
     }
 
-    @Override
-    public void pause() {
-
+    private ClientConfig getClientConfig(ConnectRecord connectRecord) {
+        ClientConfig clientConfig = new ClientConfig();
+        clientConfig.setHttpClient(httpClient);
+        clientConfig.setUrlPattern(urlPattern);
+        clientConfig.setMethod(CheckUtils.checkNull(method) ? 
connectRecord.getExtension(HttpConstant.HTTP_METHOD) : method);
+        clientConfig.setAuthType(authType);
+        
clientConfig.setHttpPathValue(connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE));
+        
clientConfig.setQueryStringParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)).toJSONString());
+        
clientConfig.setHeaderParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)).toJSONString());
+        clientConfig.setBodys(bodys);
+        clientConfig.setProxyUser(proxyUser);
+        clientConfig.setProxyPassword(proxyPassword);
+        clientConfig.setProxyType(proxyType);
+        clientConfig.setProxyPort(proxyPort);
+        clientConfig.setProxyHost(proxyHost);
+        clientConfig.setOauth2ClientId(oauth2ClientId);
+        clientConfig.setOauth2ClientSecret(oauth2ClientSecret);
+        clientConfig.setTimeout(timeout);
+        clientConfig.setOauth2HttpMethod(oauth2HttpMethod);
+        clientConfig.setOauth2Endpoint(oauth2Endpoint);
+        clientConfig.setBasicUser(basicUser);
+        clientConfig.setBasicPassword(basicPassword);
+        clientConfig.setApiKeyName(apiKeyName);
+        clientConfig.setApiKeyValue(apiKeyValue);
+        return clientConfig;
     }
 
-    @Override
-    public void resume() {
-
+    private void addHeaderMap(Map<String, String> headerMap, ClientConfig 
clientConfig) {
+        String header = clientConfig.getHeaderParameters();
+        if (StringUtils.isBlank(header)) {
+            return;
+        }
+        JSONObject jsonObject = JSONObject.parseObject(header);
+        for (Map.Entry<String, Object> entry : jsonObject.entrySet()) {
+            if (entry.getValue() instanceof JSONObject) {
+                headerMap.put(entry.getKey(), ((JSONObject) 
entry.getValue()).toJSONString());
+            } else {
+                headerMap.put(entry.getKey(), (String) entry.getValue());
+            }
+        }
     }
 
     @Override
     public void validate(KeyValue config) {
+        if 
(CheckUtils.checkNull(config.getString(HttpConstant.URL_PATTERN_CONSTANT))
+            || 
CheckUtils.checkNull(config.getString(HttpConstant.METHOD_CONSTANT))) {
+            throw new RuntimeException("http required parameter is null !");
+        }
+        final List<AuthTypeEnum> collect = 
Arrays.stream(AuthTypeEnum.values()).filter(authTypeEnum -> 
authTypeEnum.getAuthType().equals(config.getString(HttpConstant.AUTH_TYPE_CONSTANT))).collect(Collectors.toList());
+        if (collect.isEmpty()) {
+            throw new RuntimeException("authType required parameter check is 
fail !");
+        }
     }
 
     @Override
-    public void init(KeyValue config) {
-        url = config.getString(HttpConstant.URL_CONSTANT);
+    public void start(KeyValue config) {
+        urlPattern = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.URL_PATTERN_CONSTANT));
+        method = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.METHOD_CONSTANT));
+        bodys = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.BODYS_CONSTANT));
+        authType = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.AUTH_TYPE_CONSTANT));
+        basicUser = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.BASIC_USER_CONSTANT));
+        basicPassword = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.BASIC_PASSWORD_CONSTANT));
+        oauth2Endpoint = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_ENDPOINT_CONSTANT));
+        oauth2ClientId = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_CLIENTID_CONSTANT));
+        oauth2ClientSecret = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_CLIENTSECRET_CONSTANT));
+        oauth2HttpMethod = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_HTTP_METHOD_CONSTANT));
+        proxyType = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_TYPE_CONSTANT));
+        proxyHost = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_HOST_CONSTANT));

Review Comment:
   scheduledExecutorService is never initialized. The field is declared but 
only set via setScheduledExecutorService(), which is never called internally. 
In start(), scheduledExecutorService.scheduleAtFixedRate(...) will throw 
NullPointerException, preventing the task from starting entirely.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());
+                httpRequest.setHeaderMap(headerMap);
+                httpRequest.setMethod(clientConfig.getMethod());
+                httpRequest.setTimeout(clientConfig.getTimeout());
+                
httpRequest.setUrl(JsonUtils.queryStringAndPathValue(clientConfig.getUrlPattern(),
 clientConfig.getQueryStringParameters(), 
connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE)));
+                httpClient.execute(httpRequest, httpCallback);
+            }
+            boolean consumeSucceed = Boolean.FALSE;
+            try {
+                consumeSucceed = 
countDownLatch.await(DEFAULT_CONSUMER_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+            } catch (Throwable e) {
+                log.error("count down latch failed.", e);
+            }
+            if (!consumeSucceed) {
+                throw new RetriableException("Request Timeout");
+            }
+            if (httpCallback.isFailed()) {
+                throw new RetriableException(httpCallback.getMsg());
+            }
         } catch (Exception e) {

Review Comment:
   The catch(Exception e) block wraps RetriableException (thrown at lines 97 
and 100) in a new RuntimeException. This prevents the connector framework from 
recognizing retriable failures and retrying them. RetriableException should be 
caught separately and re-thrown unwrapped, or the catch block should not catch 
it.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/ApacheHttpClientImpl.java:
##########
@@ -0,0 +1,263 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import com.google.common.net.MediaType;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.http.HttpEntity;
+import org.apache.http.HttpHost;
+import org.apache.http.client.config.RequestConfig;
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpDelete;
+import org.apache.http.client.methods.HttpGet;
+import org.apache.http.client.methods.HttpHead;
+import org.apache.http.client.methods.HttpOptions;
+import org.apache.http.client.methods.HttpPatch;
+import org.apache.http.client.methods.HttpPost;
+import org.apache.http.client.methods.HttpPut;
+import org.apache.http.client.methods.HttpRequestBase;
+import org.apache.http.client.methods.HttpTrace;
+import org.apache.http.client.protocol.HttpClientContext;
+import org.apache.http.config.Registry;
+import org.apache.http.config.RegistryBuilder;
+import org.apache.http.conn.DnsResolver;
+import org.apache.http.conn.socket.ConnectionSocketFactory;
+import org.apache.http.conn.socket.PlainConnectionSocketFactory;
+import org.apache.http.conn.ssl.SSLConnectionSocketFactory;
+import org.apache.http.conn.ssl.TrustStrategy;
+import org.apache.http.entity.StringEntity;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
+import org.apache.http.protocol.HTTP;
+import org.apache.http.protocol.HttpContext;
+import org.apache.http.ssl.SSLContextBuilder;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.slf4j.MDC;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.SSLSession;
+import java.io.IOException;
+import java.io.UnsupportedEncodingException;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.Proxy;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import java.security.cert.X509Certificate;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import static 
org.apache.rocketmq.connect.http.sink.constant.HttpConstant.LOG_SIFT_TAG;
+
+
+public class ApacheHttpClientImpl implements AbstractHttpClient {
+    private static final Logger log = 
LoggerFactory.getLogger(ApacheHttpClientImpl.class);
+
+    private static ExecutorService executorServicePool = new 
ThreadPoolExecutor(200, 2000, 600, TimeUnit.SECONDS,
+            new LinkedBlockingDeque<Runnable>(1000), new 
DefaultThreadFactory("ApacheHttpClientRequestThread"));
+    private CloseableHttpClient httpClient = null;
+
+    private SocksProxyConfig socksProxyConfig;
+    private static final String SOCKS_ADDRESS_KEY = "socks.address";
+
+    @Override
+    public void init(ClientConfig config) {
+        try {
+            SSLContextBuilder sslContextBuilder = new 
SSLContextBuilder().loadTrustMaterial(null, new TrustStrategy() {
+                @Override
+                public boolean isTrusted(X509Certificate[] chain, String 
authType) {
+                    return true;
+                }
+            });
+            Registry<ConnectionSocketFactory> reg = 
RegistryBuilder.<ConnectionSocketFactory>create().register("http",
+                            new SocksPlainConnectionSocketFactory())
+                    .register("https", new 
SocksSSLConnectionSocketFactory(sslContextBuilder.build()))
+                    .build();
+            PoolingHttpClientConnectionManager connManager = new 
PoolingHttpClientConnectionManager(reg,
+                    new FakeDnsResolver());
+            connManager.setMaxTotal(400);
+            connManager.setDefaultMaxPerRoute(500);
+            httpClient = HttpClients.custom()
+                    .setConnectionManager(connManager)
+                    .build();
+
+            this.socksProxyConfig = new 
SocksProxyConfig(config.getProxyHost(), config.getProxyUser(), 
config.getProxyPassword());
+        } catch (Exception e) {
+            log.error("ApacheHttpClientImpl | init | error => ", e);
+            throw new RuntimeException(e);
+        }
+    }
+
+    @Override
+    public String execute(HttpRequest httpRequest, HttpCallback httpCallback) 
throws Exception {
+        CloseableHttpResponse response;
+        HttpRequestBase httpRequestBase = null;
+        if (httpRequest != null) {
+            httpRequestBase = extracted(httpRequest.getUrl(), 
httpRequest.getMethod(), httpRequest.getHeaderMap(), httpRequest.getBody());
+            if (StringUtils.isNotBlank(httpRequest.getTimeout())) {
+                final RequestConfig requestConfig = RequestConfig.custom().
+                        
setConnectionRequestTimeout(Integer.parseInt(httpRequest.getTimeout())).
+                        
setSocketTimeout(Integer.parseInt(httpRequest.getTimeout())).
+                        
setConnectTimeout(Integer.parseInt(httpRequest.getTimeout())).build();
+                httpRequestBase.setConfig(requestConfig);
+            }
+        }
+        HttpRequestCallable httpRequestCallable = new 
HttpRequestCallable(httpClient, httpRequestBase,
+                HttpClientContext.create(), this.socksProxyConfig, 
httpCallback, MDC.get(LOG_SIFT_TAG));
+        Future<String> submit = 
executorServicePool.submit(httpRequestCallable);
+        String result = submit.get();
+        log.info("ApacheHttpClientImpl | execute| success | result : {}", 
result);

Review Comment:
   execute() calls submit.get() which blocks indefinitely with no timeout. This 
makes the countDownLatch-based 30-second timeout in put() non-functional: if 
any HTTP request hangs and no 'timeout' config is set, submit.get() blocks 
forever and the latch await is never reached. Use submit.get(timeout, TimeUnit) 
or Future.cancel() on timeout.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/BasicAuthImpl.java:
##########
@@ -0,0 +1,31 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import com.google.common.collect.Maps;
+import com.sun.org.apache.xerces.internal.impl.dv.util.Base64;

Review Comment:
   Uses com.sun.org.apache.xercesinternal.impl.dv.util.Base64, an internal JDK 
class not part of the public API. It may not exist in non-Oracle JDKs or future 
JDK versions, causing ClassNotFoundException at runtime. Should use 
java.util.Base64.getEncoder().encodeToString() instead.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/ApacheHttpClientImpl.java:
##########
@@ -0,0 +1,263 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import com.google.common.net.MediaType;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.http.HttpEntity;
+import org.apache.http.HttpHost;
+import org.apache.http.client.config.RequestConfig;
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpDelete;
+import org.apache.http.client.methods.HttpGet;
+import org.apache.http.client.methods.HttpHead;
+import org.apache.http.client.methods.HttpOptions;
+import org.apache.http.client.methods.HttpPatch;
+import org.apache.http.client.methods.HttpPost;
+import org.apache.http.client.methods.HttpPut;
+import org.apache.http.client.methods.HttpRequestBase;
+import org.apache.http.client.methods.HttpTrace;
+import org.apache.http.client.protocol.HttpClientContext;
+import org.apache.http.config.Registry;
+import org.apache.http.config.RegistryBuilder;
+import org.apache.http.conn.DnsResolver;
+import org.apache.http.conn.socket.ConnectionSocketFactory;
+import org.apache.http.conn.socket.PlainConnectionSocketFactory;
+import org.apache.http.conn.ssl.SSLConnectionSocketFactory;
+import org.apache.http.conn.ssl.TrustStrategy;
+import org.apache.http.entity.StringEntity;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
+import org.apache.http.protocol.HTTP;
+import org.apache.http.protocol.HttpContext;
+import org.apache.http.ssl.SSLContextBuilder;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.slf4j.MDC;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.SSLSession;
+import java.io.IOException;
+import java.io.UnsupportedEncodingException;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.Proxy;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import java.security.cert.X509Certificate;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import static 
org.apache.rocketmq.connect.http.sink.constant.HttpConstant.LOG_SIFT_TAG;
+
+
+public class ApacheHttpClientImpl implements AbstractHttpClient {
+    private static final Logger log = 
LoggerFactory.getLogger(ApacheHttpClientImpl.class);
+
+    private static ExecutorService executorServicePool = new 
ThreadPoolExecutor(200, 2000, 600, TimeUnit.SECONDS,
+            new LinkedBlockingDeque<Runnable>(1000), new 
DefaultThreadFactory("ApacheHttpClientRequestThread"));
+    private CloseableHttpClient httpClient = null;
+
+    private SocksProxyConfig socksProxyConfig;
+    private static final String SOCKS_ADDRESS_KEY = "socks.address";
+
+    @Override
+    public void init(ClientConfig config) {
+        try {
+            SSLContextBuilder sslContextBuilder = new 
SSLContextBuilder().loadTrustMaterial(null, new TrustStrategy() {
+                @Override
+                public boolean isTrusted(X509Certificate[] chain, String 
authType) {
+                    return true;
+                }
+            });
+            Registry<ConnectionSocketFactory> reg = 
RegistryBuilder.<ConnectionSocketFactory>create().register("http",
+                            new SocksPlainConnectionSocketFactory())
+                    .register("https", new 
SocksSSLConnectionSocketFactory(sslContextBuilder.build()))
+                    .build();
+            PoolingHttpClientConnectionManager connManager = new 
PoolingHttpClientConnectionManager(reg,
+                    new FakeDnsResolver());
+            connManager.setMaxTotal(400);
+            connManager.setDefaultMaxPerRoute(500);
+            httpClient = HttpClients.custom()
+                    .setConnectionManager(connManager)
+                    .build();
+
+            this.socksProxyConfig = new 
SocksProxyConfig(config.getProxyHost(), config.getProxyUser(), 
config.getProxyPassword());
+        } catch (Exception e) {
+            log.error("ApacheHttpClientImpl | init | error => ", e);
+            throw new RuntimeException(e);
+        }
+    }
+
+    @Override
+    public String execute(HttpRequest httpRequest, HttpCallback httpCallback) 
throws Exception {
+        CloseableHttpResponse response;
+        HttpRequestBase httpRequestBase = null;
+        if (httpRequest != null) {
+            httpRequestBase = extracted(httpRequest.getUrl(), 
httpRequest.getMethod(), httpRequest.getHeaderMap(), httpRequest.getBody());
+            if (StringUtils.isNotBlank(httpRequest.getTimeout())) {
+                final RequestConfig requestConfig = RequestConfig.custom().
+                        
setConnectionRequestTimeout(Integer.parseInt(httpRequest.getTimeout())).
+                        
setSocketTimeout(Integer.parseInt(httpRequest.getTimeout())).
+                        
setConnectTimeout(Integer.parseInt(httpRequest.getTimeout())).build();
+                httpRequestBase.setConfig(requestConfig);
+            }
+        }
+        HttpRequestCallable httpRequestCallable = new 
HttpRequestCallable(httpClient, httpRequestBase,
+                HttpClientContext.create(), this.socksProxyConfig, 
httpCallback, MDC.get(LOG_SIFT_TAG));
+        Future<String> submit = 
executorServicePool.submit(httpRequestCallable);
+        String result = submit.get();
+        log.info("ApacheHttpClientImpl | execute| success | result : {}", 
result);
+        return result;
+    }
+
+    private HttpRequestBase extracted(String url, String method, Map<String, 
String> headerMap, String body) throws UnsupportedEncodingException {
+        switch (method) {
+            case "GET":
+                HttpGet httpGet = new HttpGet(url);
+                headerMap.forEach(httpGet::addHeader);
+                httpGet.addHeader(HTTP.CONTENT_TYPE, 
MediaType.JSON_UTF_8.toString());
+                return httpGet;
+            case "POST":
+                HttpPost httpPost = new HttpPost(url);
+                headerMap.forEach(httpPost::addHeader);
+                httpPost.addHeader(HTTP.CONTENT_TYPE, 
MediaType.JSON_UTF_8.toString());
+                if (StringUtils.isNotBlank(body)) {
+                    HttpEntity entityPot = new StringEntity(body);
+                    httpPost.setEntity(entityPot);
+                }
+                return httpPost;
+            case "DELETE":
+                HttpDelete httpDelete = new HttpDelete(url);
+                headerMap.forEach(httpDelete::addHeader);
+                httpDelete.addHeader(HTTP.CONTENT_TYPE, 
MediaType.JSON_UTF_8.toString());
+                return httpDelete;
+            case "PUT":
+                HttpPut httpPut = new HttpPut(url);
+                headerMap.forEach(httpPut::addHeader);
+                httpPut.addHeader(HTTP.CONTENT_TYPE, 
MediaType.JSON_UTF_8.toString());
+                if (StringUtils.isNotBlank(body)) {
+                    HttpEntity entityPot = new StringEntity(body);
+                    httpPut.setEntity(entityPot);
+                }
+                return httpPut;
+            case "HEAD":
+                HttpHead httpHead = new HttpHead(url);
+                headerMap.forEach(httpHead::addHeader);
+                httpHead.addHeader(HTTP.CONTENT_TYPE, 
MediaType.JSON_UTF_8.toString());
+                return httpHead;
+            case "TRACE":
+                HttpTrace httpTrace = new HttpTrace(url);
+                headerMap.forEach(httpTrace::addHeader);
+                httpTrace.addHeader(HTTP.CONTENT_TYPE, 
MediaType.JSON_UTF_8.toString());

Review Comment:
   In the extracted() switch, the TRACE case uses 'break' instead of 'return'. 
After break, execution falls through to the code after the switch that 
constructs and returns an HttpOptions object. TRACE requests will actually send 
OPTIONS requests, which is a silent correctness bug.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());
+                httpRequest.setHeaderMap(headerMap);
+                httpRequest.setMethod(clientConfig.getMethod());
+                httpRequest.setTimeout(clientConfig.getTimeout());
+                
httpRequest.setUrl(JsonUtils.queryStringAndPathValue(clientConfig.getUrlPattern(),
 clientConfig.getQueryStringParameters(), 
connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE)));
+                httpClient.execute(httpRequest, httpCallback);
+            }
+            boolean consumeSucceed = Boolean.FALSE;
+            try {
+                consumeSucceed = 
countDownLatch.await(DEFAULT_CONSUMER_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+            } catch (Throwable e) {
+                log.error("count down latch failed.", e);
+            }
+            if (!consumeSucceed) {
+                throw new RetriableException("Request Timeout");
+            }
+            if (httpCallback.isFailed()) {
+                throw new RetriableException(httpCallback.getMsg());
+            }
         } catch (Exception e) {
             log.error("HttpSinkTask | put | error => ", e);
+            throw new RuntimeException(e);
         }
     }
 
-    @Override
-    public void pause() {
-
+    private ClientConfig getClientConfig(ConnectRecord connectRecord) {
+        ClientConfig clientConfig = new ClientConfig();
+        clientConfig.setHttpClient(httpClient);
+        clientConfig.setUrlPattern(urlPattern);
+        clientConfig.setMethod(CheckUtils.checkNull(method) ? 
connectRecord.getExtension(HttpConstant.HTTP_METHOD) : method);
+        clientConfig.setAuthType(authType);
+        
clientConfig.setHttpPathValue(connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE));
+        
clientConfig.setQueryStringParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)).toJSONString());
+        
clientConfig.setHeaderParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)).toJSONString());
+        clientConfig.setBodys(bodys);
+        clientConfig.setProxyUser(proxyUser);
+        clientConfig.setProxyPassword(proxyPassword);
+        clientConfig.setProxyType(proxyType);
+        clientConfig.setProxyPort(proxyPort);
+        clientConfig.setProxyHost(proxyHost);
+        clientConfig.setOauth2ClientId(oauth2ClientId);
+        clientConfig.setOauth2ClientSecret(oauth2ClientSecret);
+        clientConfig.setTimeout(timeout);
+        clientConfig.setOauth2HttpMethod(oauth2HttpMethod);
+        clientConfig.setOauth2Endpoint(oauth2Endpoint);
+        clientConfig.setBasicUser(basicUser);
+        clientConfig.setBasicPassword(basicPassword);
+        clientConfig.setApiKeyName(apiKeyName);
+        clientConfig.setApiKeyValue(apiKeyValue);
+        return clientConfig;
     }
 
-    @Override
-    public void resume() {
-
+    private void addHeaderMap(Map<String, String> headerMap, ClientConfig 
clientConfig) {
+        String header = clientConfig.getHeaderParameters();
+        if (StringUtils.isBlank(header)) {
+            return;
+        }
+        JSONObject jsonObject = JSONObject.parseObject(header);
+        for (Map.Entry<String, Object> entry : jsonObject.entrySet()) {
+            if (entry.getValue() instanceof JSONObject) {
+                headerMap.put(entry.getKey(), ((JSONObject) 
entry.getValue()).toJSONString());
+            } else {
+                headerMap.put(entry.getKey(), (String) entry.getValue());
+            }
+        }
     }
 
     @Override
     public void validate(KeyValue config) {
+        if 
(CheckUtils.checkNull(config.getString(HttpConstant.URL_PATTERN_CONSTANT))
+            || 
CheckUtils.checkNull(config.getString(HttpConstant.METHOD_CONSTANT))) {
+            throw new RuntimeException("http required parameter is null !");
+        }
+        final List<AuthTypeEnum> collect = 
Arrays.stream(AuthTypeEnum.values()).filter(authTypeEnum -> 
authTypeEnum.getAuthType().equals(config.getString(HttpConstant.AUTH_TYPE_CONSTANT))).collect(Collectors.toList());
+        if (collect.isEmpty()) {
+            throw new RuntimeException("authType required parameter check is 
fail !");
+        }
     }
 
     @Override
-    public void init(KeyValue config) {
-        url = config.getString(HttpConstant.URL_CONSTANT);
+    public void start(KeyValue config) {
+        urlPattern = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.URL_PATTERN_CONSTANT));
+        method = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.METHOD_CONSTANT));
+        bodys = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.BODYS_CONSTANT));
+        authType = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.AUTH_TYPE_CONSTANT));
+        basicUser = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.BASIC_USER_CONSTANT));
+        basicPassword = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.BASIC_PASSWORD_CONSTANT));
+        oauth2Endpoint = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_ENDPOINT_CONSTANT));
+        oauth2ClientId = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_CLIENTID_CONSTANT));
+        oauth2ClientSecret = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_CLIENTSECRET_CONSTANT));
+        oauth2HttpMethod = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.OAUTH2_HTTP_METHOD_CONSTANT));
+        proxyType = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_TYPE_CONSTANT));
+        proxyHost = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_HOST_CONSTANT));
+        proxyPort = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_PORT_CONSTANT));
+        proxyUser = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_USER_CONSTANT));
+        proxyPassword = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.PROXY_PASSWORD_CONSTANT));
+        timeout = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.TIMEOUT_CONSTANT));
+        apiKeyName = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.API_KEY_NAME));
+        apiKeyValue = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.API_KEY_VALUE));
+        queryStringParameters = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.QUERY_STRING_PARAMETERS_CONSTANT));
+        headerParameters = 
CheckUtils.checkNullReturnDefault(config.getString(HttpConstant.HEADER_PARAMETERS_CONSTANT));
+        try {
+            httpClient = new ApacheHttpClientImpl();

Review Comment:
   stop() calls httpClient.close() without null-checking httpClient. If start() 
fails before httpClient is initialized (e.g., due to the 
scheduledExecutorService NPE), stop() will throw NullPointerException. 
Additionally, scheduledExecutorService is never shut down in stop() if it was 
set.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/HttpCallback.java:
##########
@@ -0,0 +1,57 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import org.apache.http.concurrent.FutureCallback;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.CountDownLatch;
+
+public class HttpCallback implements FutureCallback<String> {
+
+    private static final Logger log = 
LoggerFactory.getLogger(HttpCallback.class);
+
+    private CountDownLatch countDownLatch;
+
+    private boolean isFailed;
+
+    private String msg;
+
+    public HttpCallback(CountDownLatch countDownLatch) {
+        this.countDownLatch = countDownLatch;
+    }
+
+    @Override
+    public void completed(String s) {
+        countDownLatch.countDown();
+    }
+
+    public void failed(final Exception ex) {

Review Comment:
   failed() sets isFailed=true but never sets msg. When put() throws new 
RetriableException(httpCallback.getMsg()), getMsg() returns null, producing an 
exception with a null message and losing the actual error details. Should call 
setMsg(ex.getMessage()).



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/ApacheHttpClientImpl.java:
##########
@@ -0,0 +1,263 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import com.google.common.net.MediaType;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.http.HttpEntity;
+import org.apache.http.HttpHost;
+import org.apache.http.client.config.RequestConfig;
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpDelete;
+import org.apache.http.client.methods.HttpGet;
+import org.apache.http.client.methods.HttpHead;
+import org.apache.http.client.methods.HttpOptions;
+import org.apache.http.client.methods.HttpPatch;
+import org.apache.http.client.methods.HttpPost;
+import org.apache.http.client.methods.HttpPut;
+import org.apache.http.client.methods.HttpRequestBase;
+import org.apache.http.client.methods.HttpTrace;
+import org.apache.http.client.protocol.HttpClientContext;
+import org.apache.http.config.Registry;
+import org.apache.http.config.RegistryBuilder;
+import org.apache.http.conn.DnsResolver;
+import org.apache.http.conn.socket.ConnectionSocketFactory;
+import org.apache.http.conn.socket.PlainConnectionSocketFactory;
+import org.apache.http.conn.ssl.SSLConnectionSocketFactory;
+import org.apache.http.conn.ssl.TrustStrategy;
+import org.apache.http.entity.StringEntity;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
+import org.apache.http.protocol.HTTP;
+import org.apache.http.protocol.HttpContext;
+import org.apache.http.ssl.SSLContextBuilder;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.slf4j.MDC;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.SSLSession;
+import java.io.IOException;
+import java.io.UnsupportedEncodingException;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.Proxy;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import java.security.cert.X509Certificate;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import static 
org.apache.rocketmq.connect.http.sink.constant.HttpConstant.LOG_SIFT_TAG;
+
+
+public class ApacheHttpClientImpl implements AbstractHttpClient {
+    private static final Logger log = 
LoggerFactory.getLogger(ApacheHttpClientImpl.class);
+
+    private static ExecutorService executorServicePool = new 
ThreadPoolExecutor(200, 2000, 600, TimeUnit.SECONDS,
+            new LinkedBlockingDeque<Runnable>(1000), new 
DefaultThreadFactory("ApacheHttpClientRequestThread"));
+    private CloseableHttpClient httpClient = null;
+
+    private SocksProxyConfig socksProxyConfig;
+    private static final String SOCKS_ADDRESS_KEY = "socks.address";
+
+    @Override
+    public void init(ClientConfig config) {
+        try {
+            SSLContextBuilder sslContextBuilder = new 
SSLContextBuilder().loadTrustMaterial(null, new TrustStrategy() {
+                @Override
+                public boolean isTrusted(X509Certificate[] chain, String 
authType) {
+                    return true;
+                }
+            });
+            Registry<ConnectionSocketFactory> reg = 
RegistryBuilder.<ConnectionSocketFactory>create().register("http",
+                            new SocksPlainConnectionSocketFactory())
+                    .register("https", new 
SocksSSLConnectionSocketFactory(sslContextBuilder.build()))
+                    .build();
+            PoolingHttpClientConnectionManager connManager = new 
PoolingHttpClientConnectionManager(reg,
+                    new FakeDnsResolver());
+            connManager.setMaxTotal(400);
+            connManager.setDefaultMaxPerRoute(500);

Review Comment:
   Connection pool configuration is illogical: setMaxTotal(400) is less than 
setDefaultMaxPerRoute(500). Since per-route connections cannot exceed the 
total, the effective per-route limit is capped at 400. These values should be 
consistent, with maxTotal >= defaultMaxPerRoute.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());

Review Comment:
   The actual message data (connectRecord.getData()) is not included in the 
HTTP request body. The body is set from the static 'bodys' config value, 
meaning every record sends the same fixed body. This is a functional regression 
from the old behavior where connectRecord.getData().toString() was sent as the 
request payload.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());
+                httpRequest.setHeaderMap(headerMap);
+                httpRequest.setMethod(clientConfig.getMethod());
+                httpRequest.setTimeout(clientConfig.getTimeout());
+                
httpRequest.setUrl(JsonUtils.queryStringAndPathValue(clientConfig.getUrlPattern(),
 clientConfig.getQueryStringParameters(), 
connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE)));
+                httpClient.execute(httpRequest, httpCallback);
+            }
+            boolean consumeSucceed = Boolean.FALSE;
+            try {
+                consumeSucceed = 
countDownLatch.await(DEFAULT_CONSUMER_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+            } catch (Throwable e) {
+                log.error("count down latch failed.", e);
+            }
+            if (!consumeSucceed) {
+                throw new RetriableException("Request Timeout");
+            }
+            if (httpCallback.isFailed()) {
+                throw new RetriableException(httpCallback.getMsg());
+            }
         } catch (Exception e) {
             log.error("HttpSinkTask | put | error => ", e);
+            throw new RuntimeException(e);
         }
     }
 
-    @Override
-    public void pause() {
-
+    private ClientConfig getClientConfig(ConnectRecord connectRecord) {
+        ClientConfig clientConfig = new ClientConfig();
+        clientConfig.setHttpClient(httpClient);
+        clientConfig.setUrlPattern(urlPattern);
+        clientConfig.setMethod(CheckUtils.checkNull(method) ? 
connectRecord.getExtension(HttpConstant.HTTP_METHOD) : method);
+        clientConfig.setAuthType(authType);
+        
clientConfig.setHttpPathValue(connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE));
+        
clientConfig.setQueryStringParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)).toJSONString());
+        
clientConfig.setHeaderParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)).toJSONString());
+        clientConfig.setBodys(bodys);
+        clientConfig.setProxyUser(proxyUser);
+        clientConfig.setProxyPassword(proxyPassword);
+        clientConfig.setProxyType(proxyType);
+        clientConfig.setProxyPort(proxyPort);
+        clientConfig.setProxyHost(proxyHost);
+        clientConfig.setOauth2ClientId(oauth2ClientId);
+        clientConfig.setOauth2ClientSecret(oauth2ClientSecret);
+        clientConfig.setTimeout(timeout);
+        clientConfig.setOauth2HttpMethod(oauth2HttpMethod);
+        clientConfig.setOauth2Endpoint(oauth2Endpoint);
+        clientConfig.setBasicUser(basicUser);

Review Comment:
   In addHeaderMap, the else branch casts entry.getValue() to String without 
verifying the actual type. If the JSON header value is an Integer, Boolean, or 
JSONArray (not a JSONObject), this throws ClassCastException. Should use 
String.valueOf(entry.getValue()) or check all possible types.



##########
connectors/rocketmq-connect-http/README.md:
##########
@@ -15,13 +15,18 @@ mvn clean install -Dmaven.test.skip=true
 
 ```
 
http://${runtime-ip}:${runtime-port}/connectors/${rocketmq-http-sink-connector-name}
-?config={"source-rocketmq":"${runtime-ip}:${runtime-port}","source-cluster":"${broker-cluster}","connector-class":"org.apache.rocketmq.connect.http.sink.HttpSinkConnector","connect-topicname"
 : "${connect-topicname}","url":"${url}"}
+?config={"source-rocketmq":"${runtime-ip}:${runtime-port}","source-cluster":"${broker-cluster}","connector-class":"HttpSinkConnector",
+"urlPattern":"${urlPattern}","method":"${method}","queryStringParameters":"${queryStringParameters}","headerParameters":"${headerParameters}","bodys":"${bodys}","authType":"${authType}","basicUser":"${basicUser}","basicPassword":"${basicPassword}",

Review Comment:
   The JSON config example has duplicate keys: proxyPort appears 3 times and 
proxyUser appears 2 times across lines 18-19. Duplicate JSON keys cause 
undefined behavior in parsers (only the last value is typically kept). These 
should be deduplicated to match the parameter table.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());
+                httpRequest.setHeaderMap(headerMap);
+                httpRequest.setMethod(clientConfig.getMethod());
+                httpRequest.setTimeout(clientConfig.getTimeout());
+                
httpRequest.setUrl(JsonUtils.queryStringAndPathValue(clientConfig.getUrlPattern(),
 clientConfig.getQueryStringParameters(), 
connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE)));
+                httpClient.execute(httpRequest, httpCallback);
+            }
+            boolean consumeSucceed = Boolean.FALSE;
+            try {
+                consumeSucceed = 
countDownLatch.await(DEFAULT_CONSUMER_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+            } catch (Throwable e) {
+                log.error("count down latch failed.", e);
+            }
+            if (!consumeSucceed) {
+                throw new RetriableException("Request Timeout");
+            }
+            if (httpCallback.isFailed()) {
+                throw new RetriableException(httpCallback.getMsg());
+            }
         } catch (Exception e) {
             log.error("HttpSinkTask | put | error => ", e);
+            throw new RuntimeException(e);
         }
     }
 
-    @Override
-    public void pause() {
-
+    private ClientConfig getClientConfig(ConnectRecord connectRecord) {
+        ClientConfig clientConfig = new ClientConfig();
+        clientConfig.setHttpClient(httpClient);
+        clientConfig.setUrlPattern(urlPattern);
+        clientConfig.setMethod(CheckUtils.checkNull(method) ? 
connectRecord.getExtension(HttpConstant.HTTP_METHOD) : method);
+        clientConfig.setAuthType(authType);
+        
clientConfig.setHttpPathValue(connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE));
+        
clientConfig.setQueryStringParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)).toJSONString());

Review Comment:
   getClientConfig calls JsonUtils.mergeJson() twice per parameter (once to 
null-check, once to get result) and parses JSON strings 4 times total. 
Additionally, if connectRecord.getExtension() returns null, 
JSONObject.parseObject(null) may throw NPE. Extension values should be 
null-checked before parsing, and mergeJson should be called once with the 
result stored.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/HttpCallback.java:
##########
@@ -0,0 +1,57 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import org.apache.http.concurrent.FutureCallback;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.CountDownLatch;
+
+public class HttpCallback implements FutureCallback<String> {
+
+    private static final Logger log = 
LoggerFactory.getLogger(HttpCallback.class);
+
+    private CountDownLatch countDownLatch;
+
+    private boolean isFailed;

Review Comment:
   isFailed is a plain boolean accessed from multiple threads (set in callback 
threads via failed(), read in put() thread via isFailed()) without volatile or 
synchronization. The put() thread may not see the updated value due to thread 
caching, causing failed requests to be treated as successful.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/auth/ApacheHttpClientImpl.java:
##########
@@ -0,0 +1,263 @@
+package org.apache.rocketmq.connect.http.sink.auth;
+
+import com.google.common.net.MediaType;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.http.HttpEntity;
+import org.apache.http.HttpHost;
+import org.apache.http.client.config.RequestConfig;
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpDelete;
+import org.apache.http.client.methods.HttpGet;
+import org.apache.http.client.methods.HttpHead;
+import org.apache.http.client.methods.HttpOptions;
+import org.apache.http.client.methods.HttpPatch;
+import org.apache.http.client.methods.HttpPost;
+import org.apache.http.client.methods.HttpPut;
+import org.apache.http.client.methods.HttpRequestBase;
+import org.apache.http.client.methods.HttpTrace;
+import org.apache.http.client.protocol.HttpClientContext;
+import org.apache.http.config.Registry;
+import org.apache.http.config.RegistryBuilder;
+import org.apache.http.conn.DnsResolver;
+import org.apache.http.conn.socket.ConnectionSocketFactory;
+import org.apache.http.conn.socket.PlainConnectionSocketFactory;
+import org.apache.http.conn.ssl.SSLConnectionSocketFactory;
+import org.apache.http.conn.ssl.TrustStrategy;
+import org.apache.http.entity.StringEntity;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
+import org.apache.http.protocol.HTTP;
+import org.apache.http.protocol.HttpContext;
+import org.apache.http.ssl.SSLContextBuilder;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.slf4j.MDC;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.SSLSession;
+import java.io.IOException;
+import java.io.UnsupportedEncodingException;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.Proxy;
+import java.net.Socket;
+import java.net.UnknownHostException;
+import java.security.cert.X509Certificate;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import static 
org.apache.rocketmq.connect.http.sink.constant.HttpConstant.LOG_SIFT_TAG;
+
+
+public class ApacheHttpClientImpl implements AbstractHttpClient {
+    private static final Logger log = 
LoggerFactory.getLogger(ApacheHttpClientImpl.class);
+
+    private static ExecutorService executorServicePool = new 
ThreadPoolExecutor(200, 2000, 600, TimeUnit.SECONDS,
+            new LinkedBlockingDeque<Runnable>(1000), new 
DefaultThreadFactory("ApacheHttpClientRequestThread"));

Review Comment:
   executorServicePool is a static field shared across all connector instances 
but is never shut down in close(). This causes a thread leak (200 core threads) 
when connectors are stopped or recreated. The pool should be an instance field 
and shut down in close(), or use a shared lifecycle-managed pool.



##########
connectors/rocketmq-connect-http/src/main/java/org/apache/rocketmq/connect/http/sink/HttpSinkTask.java:
##########
@@ -1,61 +1,226 @@
 package org.apache.rocketmq.connect.http.sink;
 
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Maps;
 import io.openmessaging.KeyValue;
 import io.openmessaging.connector.api.component.task.sink.SinkTask;
-import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
 import io.openmessaging.connector.api.data.ConnectRecord;
 import io.openmessaging.connector.api.errors.ConnectException;
-import org.apache.rocketmq.connect.http.sink.common.OkHttpUtils;
+import io.openmessaging.connector.api.errors.RetriableException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.rocketmq.connect.http.sink.auth.AbstractHttpClient;
+import org.apache.rocketmq.connect.http.sink.auth.ApacheHttpClientImpl;
+import org.apache.rocketmq.connect.http.sink.auth.ApiKeyImpl;
+import org.apache.rocketmq.connect.http.sink.auth.BasicAuthImpl;
+import org.apache.rocketmq.connect.http.sink.auth.HttpCallback;
+import org.apache.rocketmq.connect.http.sink.auth.OAuthClientImpl;
+import org.apache.rocketmq.connect.http.sink.constant.AuthTypeEnum;
 import org.apache.rocketmq.connect.http.sink.constant.HttpConstant;
+import org.apache.rocketmq.connect.http.sink.entity.ClientConfig;
+import org.apache.rocketmq.connect.http.sink.entity.HttpRequest;
+import org.apache.rocketmq.connect.http.sink.util.CheckUtils;
+import org.apache.rocketmq.connect.http.sink.util.JsonUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
 
 public class HttpSinkTask extends SinkTask {
     private static final Logger log = 
LoggerFactory.getLogger(HttpSinkTask.class);
+    private static final int DEFAULT_CONSUMER_TIMEOUT_SECONDS = 30;
 
-    private String url;
+    protected ScheduledExecutorService scheduledExecutorService;
+    protected String urlPattern;
+    protected String method;
+    protected String queryStringParameters;
+    protected String headerParameters;
+    protected String bodys;
+    protected String authType;
+    protected String basicUser;
+    protected String basicPassword;
+    protected String oauth2Endpoint;
+    protected String oauth2ClientId;
+    protected String oauth2ClientSecret;
+    protected String oauth2HttpMethod;
+    protected String proxyType;
+    protected String proxyHost;
+    protected String proxyPort;
+    protected String proxyUser;
+    protected String proxyPassword;
+    protected String apiKeyName;
+    protected String apiKeyValue;
+    protected String timeout;
+
+    private AbstractHttpClient httpClient;
+
+    private OAuthClientImpl oAuthClient;
+
+    private BasicAuthImpl basicAuth;
+
+    private ApiKeyImpl apiKey;
 
     @Override
     public void put(List<ConnectRecord> sinkRecords) throws ConnectException {
         try {
-            sinkRecords.forEach(connectRecord -> OkHttpUtils.builder()
-                    .url(url)
-                    .addParam(HttpConstant.DATA_CONSTANT, 
connectRecord.getData().toString())
-                    .post(true)
-                    .sync());
+            CountDownLatch countDownLatch = new 
CountDownLatch(sinkRecords.size());
+            HttpCallback httpCallback = new HttpCallback(countDownLatch);
+            for (ConnectRecord connectRecord : sinkRecords) {
+                ClientConfig clientConfig = getClientConfig(connectRecord);
+                Map<String, String> headerMap = Maps.newHashMap();
+                addHeaderMap(headerMap, clientConfig);
+                if (StringUtils.isNotBlank(clientConfig.getAuthType())) {
+                    headerMap.putAll(auth(clientConfig));
+                }
+                HttpRequest httpRequest = new HttpRequest();
+                httpRequest.setBody(clientConfig.getBodys());
+                httpRequest.setHeaderMap(headerMap);
+                httpRequest.setMethod(clientConfig.getMethod());
+                httpRequest.setTimeout(clientConfig.getTimeout());
+                
httpRequest.setUrl(JsonUtils.queryStringAndPathValue(clientConfig.getUrlPattern(),
 clientConfig.getQueryStringParameters(), 
connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE)));
+                httpClient.execute(httpRequest, httpCallback);
+            }
+            boolean consumeSucceed = Boolean.FALSE;
+            try {
+                consumeSucceed = 
countDownLatch.await(DEFAULT_CONSUMER_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+            } catch (Throwable e) {
+                log.error("count down latch failed.", e);
+            }
+            if (!consumeSucceed) {
+                throw new RetriableException("Request Timeout");
+            }
+            if (httpCallback.isFailed()) {
+                throw new RetriableException(httpCallback.getMsg());
+            }
         } catch (Exception e) {
             log.error("HttpSinkTask | put | error => ", e);
+            throw new RuntimeException(e);
         }
     }
 
-    @Override
-    public void pause() {
-
+    private ClientConfig getClientConfig(ConnectRecord connectRecord) {
+        ClientConfig clientConfig = new ClientConfig();
+        clientConfig.setHttpClient(httpClient);
+        clientConfig.setUrlPattern(urlPattern);
+        clientConfig.setMethod(CheckUtils.checkNull(method) ? 
connectRecord.getExtension(HttpConstant.HTTP_METHOD) : method);
+        clientConfig.setAuthType(authType);
+        
clientConfig.setHttpPathValue(connectRecord.getExtension(HttpConstant.HTTP_PATH_VALUE));
+        
clientConfig.setQueryStringParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_QUERY_VALUE)),
 JSONObject.parseObject(queryStringParameters)).toJSONString());
+        
clientConfig.setHeaderParameters(JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)) == null ? null : 
JsonUtils.mergeJson(JSONObject.parseObject(connectRecord.getExtension(HttpConstant.HTTP_HEADER)),
 JSONObject.parseObject(headerParameters)).toJSONString());
+        clientConfig.setBodys(bodys);
+        clientConfig.setProxyUser(proxyUser);
+        clientConfig.setProxyPassword(proxyPassword);
+        clientConfig.setProxyType(proxyType);
+        clientConfig.setProxyPort(proxyPort);
+        clientConfig.setProxyHost(proxyHost);
+        clientConfig.setOauth2ClientId(oauth2ClientId);
+        clientConfig.setOauth2ClientSecret(oauth2ClientSecret);
+        clientConfig.setTimeout(timeout);
+        clientConfig.setOauth2HttpMethod(oauth2HttpMethod);
+        clientConfig.setOauth2Endpoint(oauth2Endpoint);
+        clientConfig.setBasicUser(basicUser);
+        clientConfig.setBasicPassword(basicPassword);
+        clientConfig.setApiKeyName(apiKeyName);
+        clientConfig.setApiKeyValue(apiKeyValue);
+        return clientConfig;
     }
 
-    @Override
-    public void resume() {
-
+    private void addHeaderMap(Map<String, String> headerMap, ClientConfig 
clientConfig) {
+        String header = clientConfig.getHeaderParameters();
+        if (StringUtils.isBlank(header)) {
+            return;
+        }
+        JSONObject jsonObject = JSONObject.parseObject(header);
+        for (Map.Entry<String, Object> entry : jsonObject.entrySet()) {
+            if (entry.getValue() instanceof JSONObject) {
+                headerMap.put(entry.getKey(), ((JSONObject) 
entry.getValue()).toJSONString());
+            } else {
+                headerMap.put(entry.getKey(), (String) entry.getValue());
+            }
+        }
     }

Review Comment:
   validate() always checks authType against AuthTypeEnum and throws if no 
match is found. But authType is marked optional (NO) in the README parameter 
table. If authType is not configured (null), the stream filter matches no enum 
and validation throws, contradicting the documented contract.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to