You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Connect AWS S3 Sink连接器无法读取AWS Kafka Topic数据

Troubleshooting S3 Sink Connector Failing to Consume from AWS MSK (SSL Enabled)

Let's break down what's happening here and walk through actionable fixes:

First, the key clue is that your manual KafkaConsumer works with just security.protocol=SSL, but the Connect worker's internal consumer (used by the S3 Sink) can't establish a stable connection—showing empty partitions and EOFException in debug logs. This points to a mismatch between the SSL configuration of your Connect worker and the setup that works for your standalone consumer.

Here's how to resolve this:

1. Ensure Connect Worker has complete SSL truststore configuration

Your standalone consumer might be running in an environment where the MSK broker's SSL certificate is already trusted (e.g., using the system default truststore), but the Connect worker's JVM doesn't have this trust. AWS MSK uses public certificates signed by Amazon Trust Services, so you need to explicitly configure the Connect worker to trust these roots.

Update your worker properties file with these SSL settings:

plugin.path = <plugins directory>
bootstrap.servers = <Amazon MSK服务器列表>
security.protocol = SSL
# Use the system truststore (contains Amazon's root certs) or specify a custom one
ssl.truststore.location = /usr/lib/jvm/java-1.8.0-openjdk-amazon-corretto/jre/lib/security/cacerts
ssl.truststore.password = changeit

(Adjust the truststore path to match your JDK installation; the default password for the system cacerts is changeit.)

2. Explicitly override consumer SSL settings in the connector configuration

Sometimes the Connect worker's global SSL settings don't propagate correctly to the connector's internal consumer. Force the connector to use the same SSL config that works for your manual consumer by adding these to your connector properties:

name = my-connector
connector.class = io.confluent.connect.s3.S3SinkConnector
topics = some_topic
# Override consumer-specific SSL settings to match your working standalone consumer
consumer.override.security.protocol = SSL
consumer.override.ssl.truststore.location = /usr/lib/jvm/java-1.8.0-openjdk-amazon-corretto/jre/lib/security/cacerts
consumer.override.ssl.truststore.password = changeit

3. Verify network connectivity and MSK security group rules

Even though your standalone consumer works, double-check that the machine running the Connect worker has outbound access to MSK's SSL port (default 9094). Ensure your MSK cluster's security group allows incoming traffic on port 9094 from the Connect worker's IP/security group.

Test connectivity with this command:

nc -zv <msk-broker-endpoint> 9094

4. Dig deeper into SSL handshake logs

The EOFException usually happens when the SSL handshake fails abruptly. Enable more detailed SSL logging for the Connect worker by adding this to your JVM args (in your Connect startup script):

-Djavax.net.debug=ssl:handshake:verbose

This will show you exactly where the SSL handshake is failing—like if the broker's certificate isn't trusted, or if there's a protocol mismatch.

Why this works

Your manual KafkaConsumer likely inherits the system's truststore settings automatically, but the Connect worker (especially when running as a service) might be using a different JVM or custom classpath that doesn't include the necessary root certificates. By explicitly setting the truststore in both the worker and connector configs, you ensure the internal consumer uses the same trusted certificate store that works for your standalone consumer.


内容的提问来源于stack exchange,提问作者Archie

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 13:57:29