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

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时概率较低)

排查步骤

  1. 查看Elasticsearch节点日志:Kafka Connect的警告仅提示结果,Elasticsearch的节点日志会包含具体错误详情(如字段类型冲突、映射错误),可直接定位问题字段。
  2. 提取对应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"
    
  3. 检查Elasticsearch索引映射:通过API查看目标索引的映射规则,对比审计日志字段是否符合要求:
    curl -X GET "elasticsearch_url/confluent-audit-log-events/_mapping"
    

配置调整建议

  1. 修改错误处理策略:将behavior.on.malformed.documents从WARN改为FAIL,遇到非法文档时任务会停止,避免静默丢数据,方便快速定位问题:
    "behavior.on.malformed.documents": "FAIL"
    
  2. 修复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';"
    
  3. 添加字段转换处理:若日志字段存在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"
      
  4. 手动配置索引映射:若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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 14:47:30