如何解决Kafka Connect Cassandra Sink连接器连接AWS Keyspaces失败问题?
问题分析与解决配置遗漏
你的核心问题是缺少AWS Keyspaces所需的IAM认证配置,以及部分不兼容的参数设置,以下是具体需要补充和调整的配置项:
1. 添加IAM认证相关配置
AWS Keyspaces使用IAM SigV4进行身份验证,Datastax Kafka连接器默认不会启用该认证方式,必须显式配置:
- 添加
auth.provider.class: com.datastax.oss.driver.api.core.auth.AwsSigV4AuthProvider:指定使用AWS SigV4认证提供者 - 添加
basic.auth.credentials.source: INSTANCE_METADATA:让驱动自动从EC2实例的元数据服务获取IAM角色凭证(你的EC2已配置对应Cassandra权限的IAM角色) - 可选补充
advanced.auth-provider.aws.region: us-east-1:明确指定AWS区域,避免自动识别错误
2. 调整一致性级别
AWS Keyspaces不支持QUORUM一致性级别,仅支持本地区域的一致性选项,需将consistencyLevel修改为LOCAL_QUORUM或LOCAL_ONE
3. 完善SSL强制配置
虽然你已配置信任库,但需明确开启SSL,添加ssl.enabled: true确保连接器强制使用SSL连接Keyspaces
4. 验证驱动版本兼容性
确保你使用的Datastax Kafka连接器版本对应的Java驱动≥4.10.0,因为AWS SigV4认证支持是在该版本之后引入的,旧版本会找不到AwsSigV4AuthProvider类
修改后的完整连接器配置示例
curl -i -X POST -H "Content-Type:application/json" localhost:8083/connectors -d '{ "name":"test_connector_1", "config":{ "contactPoints":"cassandra.us-east-1.amazonaws.com", "loadBalancing.localDc":"us-east-1", "port":"9142", "consistencyLevel":"LOCAL_QUORUM", "connector.class":"com.datastax.oss.kafka.sink.CassandraSinkConnector", "tasks.max":"1", "topics":"keyspaces_connector_test", "topic.keyspaces_connector_test.connector_keyspace.key_value_test.mapping":"client_id=key, message=value", "ssl.enabled":"true", "ssl.truststore.path":"/some/path/cassandra_truststore.jks", "ssl.truststore.password":"some password", "auth.provider.class":"com.datastax.oss.driver.api.core.auth.AwsSigV4AuthProvider", "basic.auth.credentials.source":"INSTANCE_METADATA", "advanced.auth-provider.aws.region":"us-east-1" } }'
额外检查:确认EC2实例的安全组允许出站访问9142端口到Keyspaces服务,安全组规则也可能阻断连接。
内容的提问来源于stack exchange,提问作者lealvcon
相关产品推荐
相关产品推荐

