Kafka Connect AWS S3 Sink连接器无法读取AWS Kafka Topic数据
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

