Learn Labs
9. Building Data Pipelines (Kafka Connect)

9.5 Example: file source → file sink

Response includes "tasks": [{"connector":"load-kafka-config","task":0}] and "type":"source".

bin/connect-distributed.sh config/connect-distributed.properties &
# "In a real production environment, you'll want at least TWO OR THREE of
#  these running to provide high availability."

Create a source connector (pipe Kafka's own config file into a topic):

echo '{"name":"load-kafka-config", "config":{"connector.class":
"FileStreamSource","file":"config/server.properties","topic":
"kafka-config-topic"}}' | curl -X POST -d @- http://localhost:8083/connectors \
  -H "Content-Type: application/json"

Response includes "tasks": [{"connector":"load-kafka-config","task":0}] and "type":"source".

Verify:

bin/kafka-console-consumer.sh --bootstrap-server=localhost:9092 \
  --topic kafka-config-topic --from-beginning

{"schema":{"type":"string","optional":false},"payload":"# Licensed to the Apache..."}
{"schema":{"type":"string","optional":false},"payload":"broker.id=0"}
...

"Note that by default, the JSON converter places a SCHEMA IN EACH RECORD. In this specific case, the schema is very simple — there is only a single column, named payload of type string, containing a single line from the file for each record."

Create the sink:

echo '{"name":"dump-kafka-config", "config":
{"connector.class":"FileStreamSink","file":"copy-of-server-properties",
"topics":"kafka-config-topic"}}' | curl -X POST -d @- \
  http://localhost:8083/connectors --header "content-Type:application/json"

The differences from the source config — note the plural:

PropertySourceSink
connector.classFileStreamSourceFileStreamSink
filethe source filethe destination file
topictopic — singulartopics — plural

► “you can write multiple topics into one file with the sink, while the source only allows writing into one topic.”

Delete:

curl -X DELETE http://localhost:8083/connectors/dump-kafka-config

⚠️ WARNING — FileStream connectors are demo-only

"This example uses FileStream connectors because they are simple and built into Kafka... These should NOT be used for actual production pipelines, as they have MANY LIMITATIONS AND NO RELIABILITY GUARANTEES. There are several alternatives if you want to ingest data from files: FilePulse Connector, FileSystem Connector, or SpoolDir."


On this page