使用Confluent Elasticsearch Sink Connector转换Kafka主题数据为结构化JSON失败求助
Hey there! Let's work through why your Kafka topic data isn't showing up as structured JSON for Elasticsearch. Based on your setup (Debezium 0.7.5, Confluent 4.1.0, Elasticsearch 6.1.0), the most common culprit is Debezium's default envelope format—here's how to fix it step by step:
1. Adjust Debezium MongoDB Source Connector Configuration
Debezium wraps MongoDB change events in an Envelope structure (with schema and payload fields) by default. To extract the actual structured document from the payload.new field, add the ExtractNewRecordState transform to your connector config:
{ "name": "mongodb-source-connector", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "mongodb.hosts": "<your-shard-cluster-hosts>", "mongodb.name": "mongo-cluster", "database.whitelist": "<your-target-db>", "collection.whitelist": "<your-target-db>.<your-target-collection>", // Add these transform settings "transforms": "extract", "transforms.extract.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.extract.drop.tombstones": "false" // Keep if you need to capture delete events } }
This transform strips away the outer envelope and pushes the actual MongoDB document to the top level of your Kafka messages.
2. Configure Elasticsearch Sink Connector for Raw JSON
Make sure your sink connector is set up to consume unstructured JSON (no schema attached):
{ "name": "elasticsearch-sink-connector", "config": { "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "tasks.max": "1", "topics": "<your-kafka-topic-name>", "key.ignore": "true", "connection.url": "http://<your-elasticsearch-host>:9200", "type.name": "_doc", // Critical settings for JSON handling "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false" } }
Disabling schemas.enable tells the sink to treat the message value as raw, structured JSON instead of expecting a schema wrapper.
3. Verify Kafka Topic Data Structure
Before checking Elasticsearch, confirm the Kafka messages are now structured. Run this command in your Confluent bin directory:
$ bin/kafka-console-consumer --bootstrap-server localhost:9092 --topic <your-kafka-topic-name> --from-beginning
You should see top-level JSON documents matching your MongoDB collection's structure, not the nested envelope format.
4. Handle Existing Topic Data (If Needed)
If you already have messages in the topic with the old envelope format, you'll need to either:
- Delete the topic and re-capture data (use
$ bin/kafka-topics --delete --topic <topic-name> --zookeeper <zookeeper-host>:2181) - Use a Kafka Streams job to reprocess old messages into the new structure
Quick Version Compatibility Check
Your stack versions are mostly compatible: Debezium 0.7.5 works with Confluent 4.1.0, and the Elasticsearch Sink for Confluent 4.1.0 supports Elasticsearch 6.x. Just ensure you're using the correct connector JARs for your Confluent version.
内容的提问来源于stack exchange,提问作者higee

