Kafka对接DynamoDB时connector处于degraded状态,求可用配置示例
问题根因说明
你的配置存在3个核心缺失项,是导致连接器状态为degraded的主要原因:
- 未开启schemaless模式支持:Confluent官方DynamoDB Sink连接器默认需要带schema的消息格式,你使用无schema JSON且关闭了schema.enable,需要额外开启schemaless配置
- 缺失AWS访问凭证配置:即使是本地DynamoDB也需要配置任意非空的access key和secret key,否则连接器鉴权逻辑会报错
- 如果你是用Docker部署Kafka Connect,
aws.dynamodb.endpoint填http://localhost:8000会指向Connect容器本身,而非宿主机的DynamoDB,需要替换为宿主机局域网IP或者Docker宿主访问别名(比如mac环境用http://host.docker.internal:8000,Linux环境用http://172.17.0.1:8000)
前置准备
- 提前在本地DynamoDB创建对应表,表名与topic名一致为
KAFKA_STOCK,设置:- 分区键(Hash Key):
companySymbol,类型为字符串 - 排序键(Sort Key):
txTime,类型根据你的数据选择,数值型或者字符串均可
- 分区键(Hash Key):
- 确认Kafka Connect集群已经安装了
io.confluent.connect.aws.dynamodb.DynamoDbSinkConnector插件
可运行连接器配置
{ "key.converter.schemas.enable": "false", "value.converter.schemas.enable": "false", "name": "dynamo-sink-connector", "connector.class": "io.confluent.connect.aws.dynamodb.DynamoDbSinkConnector", "tasks.max": "1", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "topics": "KAFKA_STOCK", "aws.dynamodb.pk.hash": "value.companySymbol", "aws.dynamodb.pk.sort": "value.txTime", "aws.dynamodb.endpoint": "http://host.docker.internal:8000", "aws.access.key.id": "fakeMyKeyId", "aws.secret.access.key": "fakeSecretAccessKey", "aws.dynamodb.region": "us-east-1", "aws.dynamodb.schemaless.enable": "true", "confluent.topic.bootstrap.servers": "broker:29092" }
注意:aws.dynamodb.endpoint的值根据你的部署环境自行调整,本地直接部署的Kafka Connect可以保留http://localhost:8000
验证步骤
- 提交上述配置到Kafka Connect,确认连接器状态变为RUNNING
- 向
KAFKA_STOCKtopic写入一条测试JSON消息:{"companySymbol": "AAPL", "txTime": 1699999999000, "price": 189.5, "tradeVolume": 234000} - 登录本地DynamoDB查询
KAFKA_STOCK表,即可看到对应数据已经同步写入
内容的提问来源于stack exchange,提问作者Vinayak Dhulipudi
相关产品推荐
相关产品推荐

