diff --git a/.env.example b/.env.example index 377f2cb..5b50dd2 100644 --- a/.env.example +++ b/.env.example @@ -18,6 +18,12 @@ MINIO_ROOT_PASSWORD=minio123 MINIO_API_PORT=9000 MINIO_CONSOLE_PORT=9001 +# MinIO bucket names +LAKE_BUCKET=lakepulse +BRONZE_PATH=lakepulse/bronze +SILVER_PATH=lakepulse/silver +GOLD_PATH=lakepulse/gold + # S3 endpoints for applications S3_ENDPOINT=http://minio:9000 S3_ACCESS_KEY=minio @@ -53,4 +59,5 @@ AIRFLOW_EMAIL=admin@lakepulse.com # APACHE KAFKA SETTINGS # ================================================================= KAFKA_BOOTSTRAP_SERVERS=broker:29092 -KAFKA_SCHEMA_REGISTRY_URL=http://schema-registry:8081 \ No newline at end of file +KAFKA_SCHEMA_REGISTRY_URL=http://schema-registry:8081 +CDC_TOPIC_PREFIX=cdc \ No newline at end of file diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml new file mode 100644 index 0000000..3f8cc04 --- /dev/null +++ b/.github/workflows/kafka-ci.yml @@ -0,0 +1,208 @@ +--- +name: Kafka CDC Pipeline CI + +permissions: + actions: read + contents: read + security-events: write + +on: + push: + branches: + - main + - feature/kafka-cdc + paths: + - 'docker/kafka-connect/**' + - 'kafka/connectors/**' + - 'docker-compose.yml' + - '.github/workflows/kafka-ci.yml' + pull_request: + branches: + - main + paths: + - 'docker/kafka-connect/**' + - 'kafka/connectors/**' + - 'docker-compose.yml' + +env: + PROJECT_NAME: lakepulse + REGISTRY: ghcr.io + IMAGE_NAME: ${{ github.repository_owner }}/lakepulse-kafka-connect + +jobs: + lint: + name: Lint Configuration Files + runs-on: ubuntu-latest + steps: + - name: Checkout repository + uses: actions/checkout@v4 + + - name: Set up .env file + run: cp .env.example .env + + - name: Setup Python + uses: actions/setup-python@v4 + with: + python-version: '3.11' + + - name: Install yamllint + run: pip install yamllint + + - name: Lint YAML files + run: | + yamllint docker-compose.yml + yamllint .github/workflows/ + + - name: Validate JSON connector configs + run: | + for file in kafka/connectors/*.json; do + echo "Validating $file" + python -m json.tool "$file" > /dev/null + done + + - name: Check Docker Compose syntax + run: docker compose config + + security: + name: Security Scan + runs-on: ubuntu-latest + steps: + - name: Checkout repository + uses: actions/checkout@v4 + + - name: Run Trivy vulnerability scanner + uses: aquasecurity/trivy-action@master + with: + scan-type: 'fs' + scan-ref: '.' + format: 'sarif' + output: 'trivy-results.sarif' + + - name: Upload Trivy scan results + uses: github/codeql-action/upload-sarif@v3 + if: always() + with: + sarif_file: 'trivy-results.sarif' + + build: + name: Build Kafka Connect Image + runs-on: ubuntu-latest + needs: [lint] + permissions: + contents: read + packages: write + + outputs: + image-digest: ${{ steps.build.outputs.digest }} + image-tag: ${{ steps.meta.outputs.tags }} + + steps: + - name: Checkout repository + uses: actions/checkout@v4 + + - name: Set up QEMU + uses: docker/setup-qemu-action@v3 + with: + platforms: linux/amd64,linux/arm64 + + - name: Set up Docker Buildx + uses: docker/setup-buildx-action@v3 + + - name: Log in to Container Registry + uses: docker/login-action@v3 + with: + registry: ${{ env.REGISTRY }} + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + + - name: Extract metadata + id: meta + uses: docker/metadata-action@v5 + with: + images: ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }} + tags: | + type=ref,event=branch + type=ref,event=pr + type=sha,prefix={{branch}}- + type=raw,value=latest,enable={{is_default_branch}} + + - name: Build and push Docker image + id: build + uses: docker/build-push-action@v6 + with: + context: ./docker/kafka-connect + file: ./docker/kafka-connect/Dockerfile + push: true + tags: ${{ steps.meta.outputs.tags }} + labels: ${{ steps.meta.outputs.labels }} + cache-from: type=gha + cache-to: type=gha,mode=max + platforms: linux/amd64,linux/arm64 + + test: + name: Integration Tests + runs-on: ubuntu-latest + needs: [build] + + steps: + - name: Checkout repository + uses: actions/checkout@v4 + + - name: Create test environment file + run: | + cp .env.example .env + + - name: Start test services + run: | + # Update docker-compose to use test image (use first tag only) + IMAGE_TAG=$(echo "${{ needs.build.outputs.image-tag }}" | head -n1) + sed -i "s|lakepulse/kafka-connect:latest|${IMAGE_TAG}|" \ + docker-compose.yml + + # Start minimal services for testing + docker compose --env-file .env up -d postgres broker \ + schema-registry minio + + # Wait for services to be ready + sleep 30 + + - name: Test Kafka Connect image + run: | + # Start Kafka Connect with test image + docker compose --env-file .env up -d kafka-connect + + # Wait for Connect to start + timeout 120 bash -c \ + 'until curl -f http://localhost:8083/; do sleep 5; done' + + - name: Validate connector plugins + run: | + # Check if required plugins are installed + curl -s http://localhost:8083/connector-plugins | \ + jq '.[] | select(.class | contains("PostgresConnector"))' + curl -s http://localhost:8083/connector-plugins | \ + jq '.[] | select(.class | contains("S3SinkConnector"))' + + - name: Test connector configuration validation + run: | + # Test Postgres source connector config validation + POSTGRES_URL="http://localhost:8083/connector-plugins" + POSTGRES_URL="${POSTGRES_URL}/io.debezium.connector.postgresql" + POSTGRES_URL="${POSTGRES_URL}.PostgresConnector/config/validate" + curl -X PUT "${POSTGRES_URL}" \ + -H "Content-Type: application/json" \ + -d @kafka/connectors/postgres-source.json | jq '.error_count' + + # Test S3 sink connector config validation + S3_URL="http://localhost:8083/connector-plugins" + S3_URL="${S3_URL}/io.confluent.connect.s3" + S3_URL="${S3_URL}.S3SinkConnector/config/validate" + curl -X PUT "${S3_URL}" \ + -H "Content-Type: application/json" \ + -d @kafka/connectors/s3-sink.json | jq '.error_count' + + - name: Cleanup test environment + if: always() + run: | + docker compose --env-file .env down -v + docker system prune -f diff --git a/Makefile b/Makefile index 87f3184..3bc5ef2 100644 --- a/Makefile +++ b/Makefile @@ -193,14 +193,28 @@ start-kafka: ## Start Kafka service @echo "$(GREEN)Starting Kafka service...$(RESET)" @docker compose up -d broker schema-registry kafka-connect kafka-ui @echo "$(GREEN)Waiting for Kafka to be ready...$(RESET)" - @sleep 10 # Wait for Kafka to initialize + @sleep 15 # Wait for Kafka to initialize @echo "$(GREEN)✓ Kafka service started$(RESET)" + @$(MAKE) register-all-connectors stop-kafka: ## Stop Kafka service @echo "$(YELLOW)Stopping Kafka service...$(RESET)" @docker compose down broker schema-registry kafka-connect kafka-ui @echo "$(GREEN)✓ Kafka service stopped$(RESET)" +register-connector: ## Register a specific connector (usage: make register-connector CONNECTOR=bronze-s3-sink) + @if [ -z "$(CONNECTOR)" ]; then \ + echo "$(RED)Please specify a connector: make register-connector CONNECTOR=bronze-s3-sink$(RESET)"; \ + exit 1; \ + fi + @echo "$(GREEN)Registering connector: $(CONNECTOR)$(RESET)" + @curl -s -X POST \ + http://localhost:8083/connectors \ + -H 'Content-Type: application/json' \ + -d @kafka/connectors/$(CONNECTOR).json | jq '.' || echo "$(RED)Failed to register $(CONNECTOR)$(RESET)" + @sleep 2 + @echo "$(GREEN)✓ Connector registered$(RESET)" + register-all-connectors: ## Register all connectors in kafka/connectors/ directory @echo "$(GREEN)Registering all connectors...$(RESET)" @for connector in kafka/connectors/*.json; do \ @@ -211,4 +225,114 @@ register-all-connectors: ## Register all connectors in kafka/connectors/ directo -d @$$connector | jq '.' || echo "$(RED)Failed to register $$connector$(RESET)"; \ sleep 2; \ done - @echo "$(GREEN)✓ All connectors registered$(RESET)" \ No newline at end of file + @echo "$(GREEN)✓ All connectors registered$(RESET)" + +list-connectors: ## List all registered connectors + @echo "$(CYAN)Registered Connectors:$(RESET)" + @curl -s http://localhost:8083/connectors | jq '.[]' || echo "$(RED)No connectors found or Kafka Connect not available$(RESET)" + +check-connector-status: ## Check status of all connectors + @echo "$(CYAN)Connector Status:$(RESET)" + @for connector in $$(curl -s http://localhost:8083/connectors | jq -r '.[]' 2>/dev/null); do \ + echo "$(GREEN)$$connector:$(RESET)"; \ + curl -s http://localhost:8083/connectors/$$connector/status | jq '.connector.state, .tasks[].state' || echo "$(RED)Failed to get status$(RESET)"; \ + echo ""; \ + done + +check-connector: ## Check specific connector status (usage: make check-connector CONNECTOR=bronze-s3-sink) + @if [ -z "$(CONNECTOR)" ]; then \ + echo "$(RED)Please specify a connector: make check-connector CONNECTOR=bronze-s3-sink$(RESET)"; \ + exit 1; \ + fi + @echo "$(CYAN)Checking connector: $(CONNECTOR)$(RESET)" + @curl -s http://localhost:8083/connectors/$(CONNECTOR)/status | jq '.' || echo "$(RED)Connector not found$(RESET)" + +restart-connector: ## Restart a specific connector (usage: make restart-connector CONNECTOR=bronze-s3-sink) + @if [ -z "$(CONNECTOR)" ]; then \ + echo "$(RED)Please specify a connector: make restart-connector CONNECTOR=bronze-s3-sink$(RESET)"; \ + exit 1; \ + fi + @echo "$(YELLOW)Restarting connector: $(CONNECTOR)$(RESET)" + @curl -s -X POST http://localhost:8083/connectors/$(CONNECTOR)/restart + @echo "$(GREEN)✓ Connector restart initiated$(RESET)" + +delete-connector: ## Delete a specific connector (usage: make delete-connector CONNECTOR=bronze-s3-sink) + @if [ -z "$(CONNECTOR)" ]; then \ + echo "$(RED)Please specify a connector: make delete-connector CONNECTOR=bronze-s3-sink$(RESET)"; \ + exit 1; \ + fi + @echo "$(YELLOW)Deleting connector: $(CONNECTOR)$(RESET)" + @curl -s -X DELETE http://localhost:8083/connectors/$(CONNECTOR) + @echo "$(GREEN)✓ Connector deleted$(RESET)" + +kafka-connect-logs: ## View Kafka Connect logs + @echo "$(CYAN)Kafka Connect Logs:$(RESET)" + @docker logs kafka-connect --tail 50 + +## ================================================================= +##@ CI/CD and Testing +## ================================================================= + +validate-configs: ## Validate all connector configurations + @echo "$(GREEN)Validating connector configurations...$(RESET)" + @for config in kafka/connectors/*.json; do \ + echo "$(CYAN)Validating $$config$(RESET)"; \ + python -m json.tool "$$config" > /dev/null && echo "$(GREEN)✓ Valid JSON$(RESET)" || echo "$(RED)✗ Invalid JSON$(RESET)"; \ + done + @echo "$(GREEN)✓ Configuration validation complete$(RESET)" + +test-connectors: ## Test connector configurations without deploying + @echo "$(GREEN)Testing connector configurations...$(RESET)" + @echo "$(CYAN)Starting test environment...$(RESET)" + @docker compose up -d broker schema-registry kafka-connect + @echo "$(CYAN)Waiting for Kafka Connect to be ready...$(RESET)" + @timeout 60 bash -c 'until curl -f http://localhost:8083/; do sleep 2; done' + @echo "$(CYAN)Testing connector configurations...$(RESET)" + @for config in kafka/connectors/*.json; do \ + connector_name=$$(basename "$$config" .json); \ + connector_class=$$(jq -r '.config."connector.class"' "$$config"); \ + echo "$(CYAN)Testing $$connector_name ($$connector_class)$(RESET)"; \ + curl -s -X PUT "http://localhost:8083/connector-plugins/$$connector_class/config/validate" \ + -H "Content-Type: application/json" \ + -d @"$$config" | jq '.error_count' | \ + (read errors; if [ "$$errors" -eq 0 ]; then echo "$(GREEN)✓ Valid$(RESET)"; else echo "$(RED)✗ $$errors errors$(RESET)"; fi); \ + done + @docker compose down broker schema-registry kafka-connect + @echo "$(GREEN)✓ Connector testing complete$(RESET)" + +health-check: ## Run comprehensive health check + @echo "$(GREEN)Running health check...$(RESET)" + @$(MAKE) check-connector-status + @echo "$(CYAN)Checking consumer lag...$(RESET)" + @for topic in $$(docker exec $$(docker compose ps -q broker) kafka-topics --bootstrap-server localhost:9092 --list | grep -E "(customers|orders|products)"); do \ + echo "$(CYAN)Topic: $$topic$(RESET)"; \ + docker exec $$(docker compose ps -q broker) kafka-consumer-groups \ + --bootstrap-server localhost:9092 \ + --group connect-s3-sink \ + --describe --topic $$topic 2>/dev/null || echo "$(YELLOW)No consumer group data$(RESET)"; \ + done + @echo "$(GREEN)✓ Health check complete$(RESET)" + +ci-setup: ## Setup for CI environment + @echo "$(GREEN)Setting up CI environment...$(RESET)" + @$(MAKE) validate-configs + @$(MAKE) build-all + @echo "$(GREEN)✓ CI setup complete$(RESET)" + +ci-test: ## Run CI tests + @echo "$(GREEN)Running CI tests...$(RESET)" + @$(MAKE) test-connectors + @echo "$(GREEN)✓ CI tests complete$(RESET)" + +## ================================================================= +##@ Development +## ================================================================= + +dev-reset: ## Reset development environment (keeps configs) + @echo "$(YELLOW)Resetting development environment...$(RESET)" + @docker compose down + @docker compose up -d postgres broker schema-registry minio kafka-connect + @echo "$(GREEN)Waiting for services...$(RESET)" + @sleep 30 + @$(MAKE) register-all-connectors + @echo "$(GREEN)✓ Development environment reset$(RESET)" \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml index a2996fb..68d5a9b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,4 +1,4 @@ -# Lightweight Setup with standalone Airflow for Local Development +--- services: postgres: @@ -141,7 +141,8 @@ services: - "9101:9101" environment: KAFKA_NODE_ID: 1 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT' + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: > + CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_PROCESS_ROLES: 'broker,controller' KAFKA_CONTROLLER_QUORUM_VOTERS: '1@broker:29093' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 @@ -152,15 +153,22 @@ services: KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_JMX_PORT: 9101 KAFKA_JMX_HOSTNAME: localhost - KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092' - KAFKA_LISTENERS: 'PLAINTEXT://broker:29092,CONTROLLER://broker:29093,PLAINTEXT_HOST://0.0.0.0:9092' + KAFKA_ADVERTISED_LISTENERS: > + PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092 + KAFKA_LISTENERS: > + PLAINTEXT://broker:29092,CONTROLLER://broker:29093,PLAINTEXT_HOST://0.0.0.0:9092 KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT' KAFKA_CONFLUENT_SCHEMA_REGISTRY_URL: ${KAFKA_SCHEMA_REGISTRY_URL} KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk' healthcheck: - test: ["CMD", "kafka-topics", "--bootstrap-server", "${KAFKA_BOOTSTRAP_SERVERS}", "--list"] + test: + - "CMD" + - "kafka-topics" + - "--bootstrap-server" + - "${KAFKA_BOOTSTRAP_SERVERS}" + - "--list" interval: 30s timeout: 10s retries: 5 @@ -174,10 +182,10 @@ services: broker: condition: service_healthy ports: - - "8061:8061" + - "8061:8081" environment: SCHEMA_REGISTRY_HOST_NAME: schema-registry - SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8061 + SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: ${KAFKA_BOOTSTRAP_SERVERS} networks: - lakepulse_network @@ -225,4 +233,4 @@ volumes: networks: lakepulse_network: - driver: bridge \ No newline at end of file + driver: bridge diff --git a/docker/spark/Dockerfile b/docker/spark/Dockerfile index e8e263e..d97fff9 100644 --- a/docker/spark/Dockerfile +++ b/docker/spark/Dockerfile @@ -25,20 +25,7 @@ RUN mkdir -p ${SPARK_HOME}/jars && \ https://repo1.maven.org/maven2/software/amazon/awssdk/bundle/2.32.10/bundle-2.32.10.jar && \ # PostgreSQL Driver wget -qO ${SPARK_HOME}/jars/postgresql-42.7.7.jar \ - https://repo1.maven.org/maven2/org/postgresql/postgresql/42.7.7/postgresql-42.7.7.jar && \ - # Kafka Dependencies - wget -qO ${SPARK_HOME}/jars/spark-sql-kafka-0-10_${SCALA_VERSION}-${SPARK_VERSION}.jar \ - https://repo1.maven.org/maven2/org/apache/spark/spark-sql-kafka-0-10_${SCALA_VERSION}/${SPARK_VERSION}/spark-sql-kafka-0-10_${SCALA_VERSION}-${SPARK_VERSION}.jar && \ - wget -qO ${SPARK_HOME}/jars/kafka-clients-${KAFKA_VERSION}.jar \ - https://repo1.maven.org/maven2/org/apache/kafka/kafka-clients/${KAFKA_VERSION}/kafka-clients-${KAFKA_VERSION}.jar && \ - wget -qO ${SPARK_HOME}/jars/spark-token-provider-kafka-0-10_${SCALA_VERSION}-${SPARK_VERSION}.jar \ - https://repo1.maven.org/maven2/org/apache/spark/spark-token-provider-kafka-0-10_${SCALA_VERSION}/${SPARK_VERSION}/spark-token-provider-kafka-0-10_${SCALA_VERSION}-${SPARK_VERSION}.jar && \ - wget -qO ${SPARK_HOME}/jars/commons-pool2-2.12.0.jar \ - https://repo1.maven.org/maven2/org/apache/commons/commons-pool2/2.12.0/commons-pool2-2.12.0.jar - # Avro Support - -RUN wget -qO ${SPARK_HOME}/jars/spark-avro_${SCALA_VERSION}-${SPARK_VERSION}.jar \ - https://repo1.maven.org/maven2/org/apache/spark/spark-avro_${SCALA_VERSION}/${SPARK_VERSION}/spark-avro_${SCALA_VERSION}-${SPARK_VERSION}.jar + https://repo1.maven.org/maven2/org/postgresql/postgresql/42.7.7/postgresql-42.7.7.jar COPY requirements.txt /requirements.txt RUN pip install --no-cache-dir -r /requirements.txt \ No newline at end of file diff --git a/kafka/connectors/postgres-source.json b/kafka/connectors/postgres-source.json index 3c3c441..d6a288a 100644 --- a/kafka/connectors/postgres-source.json +++ b/kafka/connectors/postgres-source.json @@ -1,5 +1,5 @@ { - "name": "postgres-source-connector", + "name": "postgres-source", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.hostname": "postgres", @@ -8,16 +8,22 @@ "database.password": "dbz", "database.dbname": "wideworldimporters", "database.server.name": "wideworldimporters", - "table.include.list": "application.*,purchasing.*,sales.*,warehouse.*", + "table.include.list": "application.cities,application.state_provinces,application.countries,application.people,sales.invoices,sales.invoice_line,sales.customers,sales.buying_groups,sales.customer_categories,warehouse.stock_items,warehouse.colors,warehouse.package_types", "plugin.name": "pgoutput", "publication.name": "dbz_publication", "slot.name": "debezium_slot", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.avro.AvroConverter", - "value.converter.schema.registry.url": "http://schema-registry:8061", + "value.converter.schema.registry.url": "http://schema-registry:8081", "topic.prefix": "cdc", "schema.history.internal.kafka.bootstrap.servers": "broker:29092", "schema.history.internal.kafka.topic": "schema-changes.wideworldimporters", - "snapshot.mode": "initial" + "snapshot.mode": "initial", + + "topic.creation.default.replication.factor": "1", + "topic.creation.default.partitions": "6", + "topic.creation.default.cleanup.policy": "delete", + "topic.creation.default.compression.type": "snappy", + "topic.creation.default.retention.ms": "604800000" } } \ No newline at end of file diff --git a/kafka/connectors/s3-sink.json b/kafka/connectors/s3-sink.json index ad43894..b6d498d 100644 --- a/kafka/connectors/s3-sink.json +++ b/kafka/connectors/s3-sink.json @@ -1,10 +1,12 @@ { - "name": "bronze-s3-sink", + "name": "s3-sink-1", "config": { "connector.class": "io.confluent.connect.s3.S3SinkConnector", - "tasks.max": "4", + "tasks.max": "2", "topics.regex": "cdc\\..*", + "topics.dir": "bronze", + "storage.class": "io.confluent.connect.s3.storage.S3Storage", "s3.bucket.name": "lakepulse-dev", "s3.region": "us-east-1", "s3.path.style.access.enabled": "true", @@ -12,35 +14,25 @@ "aws.access.key.id": "minio", "aws.secret.access.key": "minio123", - "partitioner.class": "io.confluent.connect.storage.partitioner.HourlyPartitioner", - "path.format": "'bronze'/'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH", - "partition.duration.ms": "3600000", "locale": "en-US", "timezone": "UTC", - "flush.size": "10000", - "rotate.interval.ms": "300000", - "rotate.schedule.interval.ms": "3600000", - + "flush.size": "100", + "rotate.interval.ms": "60000", + "rotate.schedule.interval.ms": "180000", + "s3.part.size": "5242880", "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat", - "schema.registry.url": "http://schema-registry:8081", - "parquet.codec": "snappy", - "schema.compatibility": "BACKWARD", - - "errors.tolerance": "all", - "errors.log.enable": "true", - "errors.log.include.messages": "true", - "transforms": "addTimestamp,addSource,addTopic", - "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value", - "transforms.addTimestamp.target.type": "Timestamp", - "transforms.addTimestamp.field": "ts_ms", + "partitioner.class": "io.confluent.connect.storage.partitioner.DailyPartitioner", + "path.format": "'bronze'/'topic'/'year'=YYYY/'month'=MM/'day'=dd", - "transforms.addSource.type": "org.apache.kafka.connect.transforms.InsertField$Value", - "transforms.addSource.static.field": "ingestion_timestamp", - "transforms.addSource.static.value": "${timestamp()}", + "consumer.override.max.poll.interval.ms": "1800000", + "value.converter.schemas.enable": "true", - "transforms.addTopic.type": "org.apache.kafka.connect.transforms.InsertField$Value", - "transforms.addTopic.topic.field": "kafka_topic" + "transforms": "insertTimestamp,insertTopic", + "transforms.insertTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value", + "transforms.insertTimestamp.timestamp.field": "ingestion_timestamp", + "transforms.insertTopic.type": "org.apache.kafka.connect.transforms.InsertField$Value", + "transforms.insertTopic.topic.field": "source_topic" } } \ No newline at end of file