CentralMesh.io

Kafka Connect
AdSense Banner (728x90)

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

json
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:

text
1s3://bucket/topics/orders/year=2024/month=01/day=15/orders+0+0000000000.parquet

Field-based partitioning uses message fields:

text
1s3://bucket/topics/orders/region=us-east/orders+0+0000000000.parquet

FILE 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.