# Kafka to S3 Pipeline
# Consume messages from Kafka and write them to S3 as batched JSON files
#
# Source: https://docs.expanso.io/examples/kafka-to-s3
#
# Usage:
#   curl -o config.yaml https://docs.expanso.io/examples/kafka-to-s3.yaml
#   expanso-edge run -f config.yaml

input:
  kafka:
    addresses:
      - localhost:9092
    topics:
      - events
    consumer_group: expanso-s3-archiver
    start_from_oldest: true

pipeline:
  processors:
    # Parse JSON payload
    - mapping: |
        root = this.parse_json().catch(deleted())

    # Add Kafka metadata for traceability
    - mapping: |
        root.kafka_topic = @kafka_topic
        root.kafka_partition = @kafka_partition
        root.kafka_offset = @kafka_offset
        root.processed_at = now()

    # Batch messages into JSON array
    - archive:
        format: json_array

output:
  aws_s3:
    bucket: my-data-lake
    path: events/${!now().ts_format("2006/01/02/15")}/batch-${!timestamp_unix_nano()}.json
    content_type: application/json
    region: us-east-1
    batching:
      count: 100
      period: 30s
