6.5 S3 Sink Connector
Land partitioned Kafka data in Amazon S3 for data lake and analytics workloads.
S3 Sink Connector
Summary
The S3 Sink Connector writes Kafka messages to S3 in various formats including JSON, Avro, and Parquet. It handles partitioning, file rotation, and S3 uploads automatically.
USE CASES
Data Lake Storage: Archive Kafka data for long-term retention and analytics.
Analytics Pipelines: Write data to S3 for processing with Spark, Athena, or Redshift Spectrum.
Compliance and Audit: Store immutable records in S3 for compliance requirements.
Backup: Create backup copies of critical Kafka data.
CONNECTOR CONFIGURATION
1{
2 "name": "s3-sink",
3 "config": {
4 "connector.class": "io.confluent.connect.s3.S3SinkConnector",
5 "tasks.max": "1",
6 "topics": "orders,transactions",
7 "s3.bucket.name": "my-data-lake",
8 "s3.region": "us-east-1",
9 "flush.size": "1000",
10 "rotate.interval.ms": "3600000",
11 "rotate.schedule.interval.ms": "3600000",
12 "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
13 "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
14 "path.format": "'year'=YYYY/'month'=MM/'day'=dd",
15 "partition.duration.ms": "86400000",
16 "timestamp.extractor": "Record"
17 }
18}FILE FORMATS
The connector supports multiple output formats:
JSON: Human-readable, easy to process. Larger file sizes.
Avro: Compact binary format with embedded schema. Good balance of size and compatibility.
Parquet: Columnar format optimized for analytics. Best compression and query performance.
PARTITIONING STRATEGIES
Time-based partitioning organizes data by date/time:
1s3://bucket/topics/orders/year=2024/month=01/day=15/orders+0+0000000000.parquetField-based partitioning uses message fields:
1s3://bucket/topics/orders/region=us-east/orders+0+0000000000.parquetFILE ROTATION
Files rotate based on:
flush.size: Number of records per file.
rotate.interval.ms: Time-based rotation.
schema.compatibility: Rotate on schema changes.
Key ideas
S3 Sink Connector streams Kafka to S3 for data lakes, analytics, and archival.
Supports JSON, Avro, and Parquet formats with configurable compression.
Time-based and field-based partitioning organize data efficiently.
Automatic file rotation based on size, time, or schema changes.
Enables analytics with Athena, Spectrum, Spark without custom code.