Getting data from Kafka into S3 sounds like a small job. Read some records, write some files. I work on the S3 Sink connector at Confluent, and most of the questions I see about it aren’t about whether it works. They’re about a handful of settings that decide how many files you get, how big they are, and how your bucket is laid out.
This post covers why you probably shouldn’t write this yourself, what Kafka Connect gives you, the S3 Sink settings worth understanding, and a small demo you can run in a few minutes.
Why not just write a consumer?
Here’s the version everyone writes first:
while (true) {
records = consumer.poll(...);
s3.putObject(bucket, key, serialize(records));
}
About ten lines. It looks done. Now the questions start:
- Crash recovery. If the process dies after the upload but before committing offsets, do you write duplicates? If it commits first, do you lose data?
- Offset management. Which offsets map to which S3 object?
- Parallelism. How do you spread partitions across several machines, and rebalance when one dies?
- Monitoring. How do you know it’s stuck?
- The next integration. When someone needs the same data in a search index next month, you write another one from scratch.
Each of these is real engineering work, and none of it is your actual product.
What Kafka Connect gives you
Kafka Connect is a framework that ships with Apache Kafka for moving data in and out of it. Source connectors bring data in from systems like Postgres, MySQL or MongoDB. Sink connectors push data out to systems like S3, Elasticsearch or Snowflake.
You don’t write Java. You write a JSON config, and the framework handles the rest:
- Horizontal scaling. Add workers and tasks spread across them.
- Fault tolerance. If a worker dies, its tasks move to the surviving workers.
- A big ecosystem. There are pre-built connectors for most common systems, so the one you need probably exists already.
- A REST API. Create, check, pause, resume and delete connectors with HTTP calls. No deployment needed to add a pipeline.
- Built-in error handling, including dead letter queues.
The S3 Sink connector
The S3 Sink connector reads from Kafka topics and writes objects to an S3 bucket. In short:
- Output formats: JSON, Avro, Parquet and raw bytes.
- Pluggable partitioners that control the folder layout in your bucket.
- Exactly-once delivery when you use a deterministic partitioner.
- Dead letter queue support for bad records.
- Works with S3-compatible storage, not only AWS.
- Can optionally store record keys and headers alongside values.
The config, setting by setting
Here’s a config shaped like one you’d use for analytics:
{
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "1",
"topics": "orders",
"s3.bucket.name": "my-data-lake",
"s3.region": "ap-south-1",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
"flush.size": "1000",
"rotate.schedule.interval.ms": "3600000",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"partition.duration.ms": "3600000",
"path.format": "'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH",
"locale": "en-US",
"timezone": "Asia/Kolkata"
}
This assumes the records are Avro with Schema Registry, set through the worker’s or connector’s value.converter. connector.class, topics, s3.bucket.name and s3.region are what they look like. The interesting ones are below.
flush.size
flush.size is the number of records written to a file before the connector commits it to S3.
This is the setting I’d look at first, because it’s a trade-off with real cost:
- Too small, and you create lots of tiny objects. Every object is a PUT request, and S3 bills per request. Lots of small files are also slower to query later.
- Too large, and records sit in the connector for longer before they show up in S3.
1000 is a reasonable place to start. Tune from there based on your record size and how fresh the data needs to be.
rotate.interval.ms vs rotate.schedule.interval.ms
These two names look almost identical and behave differently. This is the one that confuses everyone.
rotate.interval.msrotates based on record time, using the timestamps from the partitioner’s timestamp extractor. A file is closed once the records in it span the interval. It’s evaluated as records arrive, so if a partition goes quiet, its last file can stay open until the next record shows up.rotate.schedule.interval.msrotates based on the wall clock. Every interval, it commits whatever it has, whether or not new records arrived. It needstimezoneset.
If your records arrive late or in bursts, the two behave very differently. The scheduled version gives you predictable files, like one per hour. The trade-off is that wall-clock rotation isn’t deterministic, so it gives up the exactly-once guarantee. The connector docs spell out how these interact with flush.size. Whichever limit is hit first triggers the commit.
format.class
JsonFormat: easy to read, good for debugging.AvroFormat: compact, carries its schema.ParquetFormat: columnar and compressed, with codecs like snappy, gzip and zstd. For analytics, Parquet is almost always the right choice. Query engines like Athena read only the columns a query needs.ByteArrayFormat: writes the raw serialized bytes as-is.
One catch with Parquet: it needs records with a schema. In practice that means Avro, Protobuf or JSON Schema with Schema Registry, or JSON with embedded schemas. Plain schemaless JSON won’t work with Parquet.
partitioner.class
The partitioner decides the folder layout in S3. By default, objects go under a topics/ prefix, which you can change with topics.dir. File names follow <topic>+<kafka partition>+<start offset>.<format>.
DefaultPartitioner mirrors Kafka’s own layout:
topics/orders/partition=0/orders+0+0000000000.json
topics/orders/partition=1/orders+1+0000000000.json
Simple, but it says nothing about time or content, so query engines can’t skip any of it.
FieldPartitioner uses a field from the record, set with partition.field.name:
topics/orders/city=Hyderabad/orders+0+0000000000.json
topics/orders/city=Mumbai/orders+0+0000000003.json
Good when your queries always filter on that field.
TimeBasedPartitioner uses path.format and partition.duration.ms to write Hive-style time paths:
topics/orders/year=2026/month=04/day=11/hour=10/orders+0+0000000000.parquet
For a data lake, this is the one I’d reach for. When you query with something like Athena and filter by date, it only scans the matching folders instead of the whole bucket. DailyPartitioner and HourlyPartitioner are presets of the same idea. If none of these fit, you can write your own by implementing the Partitioner interface.
IAM permissions
Get these right before you start, or you’ll spend a while staring at “Access Denied”. The connector’s IAM identity needs:
s3:PutObjects3:GetObjects3:AbortMultipartUploads3:PutObjectTaggings3:ListAllMyBucketss3:ListBuckets3:GetBucketLocation
Check the connector docs for the current list for your version. For credentials, prefer an IAM role or the default AWS credential provider chain over putting access keys in the connector config. The demo below uses keys only because it runs on a laptop.
Dead letter queues
A single malformed record shouldn’t stop your whole pipeline. Connect lets you route records that fail to a dead letter queue topic instead:
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq-s3-sink",
"errors.deadletterqueue.context.headers.enable": "true"
The headers option adds the error context to each DLQ record, which makes it much easier to work out what went wrong. Then actually monitor the DLQ topic. A DLQ nobody watches just turns loud failures into quiet ones.
schema.compatibility
When records carry schemas, the S3 Sink has its own schema.compatibility setting. It controls what the connector does when the schema of incoming records changes:
NONE: every schema change closes the current file and starts a new one. Safe, but frequent schema changes mean lots of small files.BACKWARD,FORWARD,FULL: the connector checks new schemas against the current one and can keep writing compatible records into the same file. Records with an incompatible schema cause an error instead of silently producing a mismatched file.
If your producers use Schema Registry, set this to match the compatibility rule you use there. The demo below uses schemaless JSON, so it’s set to NONE.
The demo, step by step
All you need is a running Kafka broker and a Kafka Connect worker with the S3 Sink connector installed. Use whatever you already have: a local Kafka install, Confluent Platform or a cluster at work. The steps below only talk to Connect’s REST API on port 8083, so they’re the same either way.
What you need
- Kafka and a Kafka Connect worker, with the S3 Sink connector installed (
confluent-hub install confluentinc/kafka-connect-s3:latest) - The AWS CLI, signed in to an account where you can create a bucket
curlandjqAWS_REGION,AWS_ACCESS_KEY_IDandAWS_SECRET_ACCESS_KEYset in your shell
Run it
1. Check the S3 Sink connector is installed.
curl -s http://localhost:8083/connector-plugins | jq '.[].class' | grep -i S3
2. Create the bucket.
aws s3api create-bucket --bucket awsughyd-demo-ashfaq \
--region $AWS_REGION \
--create-bucket-configuration LocationConstraint=$AWS_REGION
3. Create the connector with one REST call. Using PUT on /connectors/<name>/config creates the connector if it doesn’t exist and updates it if it does, so it’s safe to re-run.
curl -s -X PUT http://localhost:8083/connectors/s3-sink-demo/config \
-H "Content-Type: application/json" \
-d '{
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "1",
"topics": "orders",
"s3.region": "'$AWS_REGION'",
"s3.bucket.name": "awsughyd-demo-ashfaq",
"aws.access.key.id": "'$AWS_ACCESS_KEY_ID'",
"aws.secret.access.key": "'$AWS_SECRET_ACCESS_KEY'",
"topics.dir": "demo",
"flush.size": "3",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"partition.duration.ms": "3600000",
"rotate.interval.ms": "3600000",
"path.format": "'\''year'\''=YYYY/'\''month'\''=MM/'\''day'\''=dd/'\''hour'\''=HH",
"locale": "en-US",
"timezone": "Asia/Kolkata",
"timestamp.extractor": "Wallclock",
"format.class": "io.confluent.connect.s3.format.json.JsonFormat",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"schema.compatibility": "NONE"
}' | jq .
The access keys come from your shell, so they never sit in a file. That’s fine on a laptop; anywhere else, give Connect an IAM role and drop those two lines. Note flush.size is set to 3, which is far too small for real use, but it makes files appear quickly on stage.
4. Check it’s running.
curl -s http://localhost:8083/connectors/s3-sink-demo/status \
| jq '{name: .name, state: .connector.state, tasks: [.tasks[].state]}'
5. Produce nine order events. Each line is key|value:
kafka-console-producer --bootstrap-server localhost:9092 \
--topic orders --property "parse.key=true" --property "key.separator=|" << 'EOF'
order-001|{"order_id":"order-001","customer":"Priya Sharma","item":"Echo Dot 5th Gen","quantity":2,"price":4499,"city":"Hyderabad","timestamp":"2026-04-11T10:15:00Z"}
order-002|{"order_id":"order-002","customer":"Rahul Verma","item":"Fire TV Stick 4K","quantity":1,"price":3999,"city":"Bangalore","timestamp":"2026-04-11T10:15:05Z"}
...
EOF
6. Look at the bucket.
aws s3 ls s3://awsughyd-demo-ashfaq/demo/ --recursive
Nine records with flush.size of 3 gives three objects, under a path like demo/orders/year=2026/month=04/day=11/hour=10/. Download one and you’ll see three JSON order records, one per line.
7. Clean up.
curl -s -X DELETE http://localhost:8083/connectors/s3-sink-demo
aws s3 rm s3://awsughyd-demo-ashfaq/demo/ --recursive
No consumer code, no offset bookkeeping, no retry logic. One config, one REST call.
Takeaways
- Writing your own Kafka-to-S3 consumer is easy to start and hard to get right. Kafka Connect already solves crash recovery, offsets, scaling and monitoring.
flush.sizeis a trade-off between PUT requests and freshness. Start around 1000 and tune.rotate.interval.msfollows record time.rotate.schedule.interval.msfollows the wall clock. Know which one you’re using.- Use Parquet and the TimeBasedPartitioner for anything you’ll query, and remember Parquet needs schemas.
- Set up IAM permissions and a dead letter queue before you go to production, not after.
I gave this as a talk at AWS User Group Hyderabad in April 2026. See the talk page.

