Elasticsearch Sink Connector运行时出现BAD_REQUEST警告求助
问题排查:Elasticsearch Sink Connector 非法文档警告
问题现象
使用Elasticsearch Sink Connector同步Confluent审计日志到Elasticsearch时,大部分数据可正常同步,但Kafka Connect日志持续出现以下警告:
[2023-03-17 15:47:28,630] WARN [elastic-sink|task-0] Encountered an illegal document error -> Response status: 'BAD_REQUEST',
2023-03-17 21:17:28 Index: 'confluent-audit-log-events',
2023-03-17 21:17:28 Document Id: 'confluent-audit-log-events+1+36008'.
2023-03-17 21:17:28 Ignoring and will not index record. (io.confluent.connect.elasticsearch.ElasticsearchClient:649)
警告影响
该警告直接导致对应ID的单条审计日志被跳过,无法写入Elasticsearch,长期积累会丢失部分日志数据,破坏审计数据的完整性。
可能原因
BAD_REQUEST是Elasticsearch返回的请求拒绝状态,常见触发原因包括:
- 日志字段类型与Elasticsearch索引映射不匹配(例如字符串写入数字类型字段)
- 文档包含Elasticsearch不允许的字段名(如带
.、@或空格的字段) - 字段长度超过Elasticsearch的索引限制
- 文档结构违反索引的动态映射规则
- 文档ID冲突(配置
write.method: UPSERT时概率较低)
排查步骤
- 查看Elasticsearch节点日志:Kafka Connect的警告仅提示结果,Elasticsearch的节点日志会包含具体错误详情(如字段类型冲突、映射错误),可直接定位问题字段。
- 提取对应ID的Kafka消息:使用
kafka-console-consumer工具获取目标ID的消息内容,检查是否存在异常字段或格式:kafka-console-consumer --bootstrap-server bootstrap --topic confluent-audit-log-events --property print.key=true --property key.separator=":" --from-beginning | grep "confluent-audit-log-events+1+36008" - 检查Elasticsearch索引映射:通过API查看目标索引的映射规则,对比审计日志字段是否符合要求:
curl -X GET "elasticsearch_url/confluent-audit-log-events/_mapping"
配置调整建议
- 修改错误处理策略:将
behavior.on.malformed.documents从WARN改为FAIL,遇到非法文档时任务会停止,避免静默丢数据,方便快速定位问题:"behavior.on.malformed.documents": "FAIL" - 修复Consumer认证配置:当前Connector配置中
consumer.override.sasl.jaas.config的用户名密码为空,若Kafka集群需要认证,需补充正确的账号信息:"consumer.override.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username='your_user' password='your_pass';" - 添加字段转换处理:若日志字段存在Elasticsearch不兼容的情况,可通过Kafka Connect Transform重命名或转换字段:
- 示例:重命名包含
.的字段"transforms": "renameFields", "transforms.renameFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.renameFields.renames": "old.field.name:new_field_name"
- 示例:重命名包含
- 手动配置索引映射:若Elasticsearch动态映射导致类型冲突,可提前创建目标索引并指定各字段的正确类型,避免自动推断错误。
附用户原始配置
Kafka Connect Docker Compose配置
--- version: '3.5' services: connect: container_name: kafka-connect image: confluentinc/cp-kafka-connect:7.2.4 ports: - "8083:8083" - "5005:5005" volumes: - ./data:/data environment: CONNECT_BOOTSTRAP_SERVERS: "bootstrap" CONNECT_GROUP_ID: "connect-7.3.1" CONNECT_CONFIG_STORAGE_TOPIC: connect-configs-7.3.1 CONNECT_OFFSET_STORAGE_TOPIC: connect-offsets-7.3.1 CONNECT_STATUS_STORAGE_TOPIC: connect-status-7.3.1 CONNECT_REPLICATION_FACTOR: 3 CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 3 CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 3 CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 3 CONNECT_KEY_CONVERTER: "org.apache.kafka.connect.storage.StringConverter" CONNECT_VALUE_CONVERTER: "io.confluent.connect.avro.AvroConverter" #CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE: "true" CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: "url" #CONNECT_VALUE_CONVERTER_BASIC_AUTH_CREDENTIALS_SOURCE: "USER_INFO" CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_BASIC_AUTH_USER_INFO: "username:secret" CONNECT_INTERNAL_KEY_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_INTERNAL_VALUE_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_REST_ADVERTISED_HOST_NAME: "connect" CONNECT_PLUGIN_PATH: /data CONNECT_LOG4J_ROOT_LOGLEVEL: INFO CONNECT_LOG4J_LOGGERS: org.reflections=ERROR #CLASSPATH: /usr/share/java/monitoring-interceptors/monitoring-interceptors-7.3.1.jar CONNECT_CONNECTOR_CLIENT_CONFIG_OVERRIDE_POLICY: All # Connect worker CONNECT_SECURITY_PROTOCOL: SASL_SSL CONNECT_SASL_JAAS_CONFIG: "org.apache.kafka.common.security.plain.PlainLoginModule required username='user' password='password';" CONNECT_SASL_MECHANISM: PLAIN CONNECT_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM: "HTTPS" # Connect consumer CONNECT_CONSUMER_SECURITY_PROTOCOL: SASL_SSL CONNECT_CONSUMER_SASL_JAAS_CONFIG: "org.apache.kafka.common.security.plain.PlainLoginModule required username="org.apache.kafka.common.security.plain.PlainLoginModule required username='user' password='password';" CONNECT_CONSUMER_SASL_MECHANISM: PLAIN CONNECT_TOPIC_CREATION_ENABLE: "true" CONNECT_LOG4J_APPENDER_STDOUT_LAYOUT_CONVERSIONPATTERN: "[%d] %p %X{connector.context}%m (%c:%L)%n"
Elasticsearch Sink Connector配置
curl -X PUT \ -H "Content-Type: application/json" \ --data '{ "tasks.max": "1", "connector.class": "ElasticsearchSinkConnector", "topics": "confluent-audit-log-events", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "consumer.override.bootstrap.servers": "bootstrap", "consumer.override.sasl.mechanism": "PLAIN", "consumer.override.security.protocol": "SASL_SSL", "consumer.override.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"\" password=\"\";", "consumer.override.client.dns.lookup": "use_all_dns_ips", "schema.ignore": "true", "key.ignore": "true", "connection.url": "elasticsearch url", "connection.username": "user", "connection.password": "password", "errors.log.enable": true, "errors.log.include.messages": true, "write.method": "UPSERT", "max.retries": "1", "retry.backoff.ms": "1000", "batch.size": "100", "behavior.on.malformed.documents": "WARN", "behavior.on.null.values": "IGNORE", "linger.ms": "10000" }' \ http://localhost:8083/connectors/elastic-sink/config | jq .
内容的提问来源于stack exchange,提问作者ravvi
相关产品推荐
相关产品推荐

