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

Kafka主题newtest关联FileStreamSink连接器任务配置咨询

Kafka Connect FileStreamSink Configuration Validation & Troubleshooting

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: FileStreamSink is the official, maintained sink connector for writing to files—perfect for simple testing like your use case.
  • file: You’ve set this to connector-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 your newtest topic—this tells the connector which topic to consume from.
  • Converters: Using StringConverter with schemas.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 has name listed 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:

  1. 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/status
    

    Look for "state": "RUNNING" in the response for both the connector and its task.

  2. Validate the Output File
    Navigate to the Connect worker’s working directory (or the absolute path you specified) and open connector-test. You should see all three messages in order (assuming your newtest topic uses a single partition—if multiple partitions, order depends on how the connector assigned partitions to tasks).

  3. Confirm Topic Messages Exist
    If the file is empty, first verify the newtest topic actually has your messages using the console consumer:

    kafka-console-consumer.sh --bootstrap-server <your-kafka-broker>:9092 --topic newtest --from-beginning
    

    This 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:
    curl -X GET http://<your-connect-worker-host>:<connect-port>/connectors/connect-test-9/config
    
    Ensure topics is set to newtest and 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:16:37