Kafka主题newtest关联FileStreamSink连接器任务配置咨询
Let’s walk through verifying your setup and making sure your newtest topic messages end up in the connector-test file as expected.
Configuration Breakdown & Validation
First, let’s review your connector config (I’ll fill in the truncated key.converter.sch part assuming it’s key.converter.schemas.enable":"false" to match your value converter setup):
{ "name":"connect-test-9", "config": { "connector.class":"FileStreamSink", "file":"connector-test", "topics":"newtest", "name":"connect-test-9", "value.converter":"org.apache.kafka.connect.storage.StringConverter", "value.converter.schemas.enable":"false", "key.converter":"org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable":"false" } }
Here’s how each critical parameter checks out:
connector.class:FileStreamSinkis the official, maintained sink connector for writing to files—perfect for simple testing like your use case.file: You’ve set this toconnector-test. Note that this is a relative path by default, so it will be created in the working directory of your Kafka Connect worker process (not where you submitted the config from). If you want a specific location, use an absolute path (e.g.,/tmp/connector-test).topics: Correctly mapped to yournewtesttopic—this tells the connector which topic to consume from.- Converters: Using
StringConverterwithschemas.enable":"false"is ideal here because your messages are plain strings (Hello, etc.). This avoids unnecessary schema serialization overhead and matches your message format exactly. - Duplicate
name: The config hasnamelisted twice, but this is harmless—Kafka Connect will just use the top-level or config-level value consistently. It’s cleaner to remove one, but it won’t break functionality.
Expected Behavior
If everything is working correctly, your three messages from newtest should be written to the connector-test file, each on a separate line:
Hello Is anybody out there Can you hear me
Step-by-Step Verification
Follow these steps to confirm your setup is functioning:
Check Connector Status
Run this API call to ensure the connector and its tasks are running:curl -X GET http://<your-connect-worker-host>:<connect-port>/connectors/connect-test-9/statusLook for
"state": "RUNNING"in the response for both the connector and its task.Validate the Output File
Navigate to the Connect worker’s working directory (or the absolute path you specified) and openconnector-test. You should see all three messages in order (assuming yournewtesttopic uses a single partition—if multiple partitions, order depends on how the connector assigned partitions to tasks).Confirm Topic Messages Exist
If the file is empty, first verify thenewtesttopic actually has your messages using the console consumer:kafka-console-consumer.sh --bootstrap-server <your-kafka-broker>:9092 --topic newtest --from-beginningThis should print your three messages to the terminal.
Common Troubleshooting Tips
- File not created? Check Connect worker logs for permission errors (the worker process needs write access to the target directory) or path issues (if using an absolute path, ensure the directory exists).
- Messages not appearing? Double-check the connector’s active config with:
Ensurecurl -X GET http://<your-connect-worker-host>:<connect-port>/connectors/connect-test-9/configtopicsis set tonewtestand converters are configured correctly. Also, verify the Connect worker can reach your Kafka brokers. - Serialization errors? Since you’ve disabled schemas and are using
StringConverter, this shouldn’t happen—but if it does, confirm your producer is sending plain string values (not binary data or structured schemas).
内容的提问来源于stack exchange,提问作者Shlomo

