From 2c5a692e307d6482f22ba311ebf48374c62ce2a2 Mon Sep 17 00:00:00 2001 From: Shadow Date: Sat, 9 Aug 2025 21:07:52 +0700 Subject: [PATCH 1/9] feat: add kafka ci workflow, fix kafka performance --- .github/workflows/kafka-ci.yml | 209 ++++++++++++++++++++++++++ Makefile | 128 +++++++++++++++- docker-compose.yml | 4 +- docker/spark/Dockerfile | 15 +- kafka/connectors/postgres-source.json | 14 +- kafka/connectors/s3-sink.json | 42 +++--- 6 files changed, 365 insertions(+), 47 deletions(-) create mode 100644 .github/workflows/kafka-ci.yml diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml new file mode 100644 index 0000000..5ef5487 --- /dev/null +++ b/.github/workflows/kafka-ci.yml @@ -0,0 +1,209 @@ +name: Kafka CDC Pipeline CI + +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: 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] + services: + postgres: + image: postgres:15-alpine + env: + POSTGRES_USER: ${{ secrets.POSTGRES_USER }} + POSTGRES_PASSWORD: ${{ secrets.POSTGRES_PASSWORD }} + POSTGRES_DB: ${{ secrets.POSTGRES_DB }} + options: >- + --health-cmd pg_isready + --health-interval 10s + --health-timeout 5s + --health-retries 5 + ports: + - 5432:5432 + + steps: + - name: Checkout repository + uses: actions/checkout@v4 + + - name: Create test environment file + run: | + cat > .env.test << EOF + POSTGRES_USER=${{ secrets.POSTGRES_USER }} + POSTGRES_PASSWORD=${{ secrets.POSTGRES_PASSWORD }} + POSTGRES_DB=${{ secrets.POSTGRES_DB }} + POSTGRES_PORT=${{ secrets.POSTGRES_PORT }} + MINIO_ROOT_USER=${{ secrets.MINIO_ROOT_USER }} + MINIO_ROOT_PASSWORD=${{ secrets.MINIO_ROOT_PASSWORD }} + KAFKA_BOOTSTRAP_SERVERS=${{ secrets.KAFKA_BOOTSTRAP_SERVERS }} + EOF + + - name: Start test services + run: | + # Update docker-compose to use test image + sed -i 's|lakepulse/kafka-connect:latest|${{ needs.build.outputs.image-tag }}|' docker-compose.yml + + # Start minimal services for testing + docker compose --env-file .env.test up -d 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.test 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 + curl -X PUT http://localhost:8083/connector-plugins/io.debezium.connector.postgresql.PostgresConnector/config/validate \ + -H "Content-Type: application/json" \ + -d @kafka/connectors/postgres-source.json | jq '.error_count' + + # Test S3 sink connector config validation + curl -X PUT http://localhost:8083/connector-plugins/io.confluent.connect.s3.S3SinkConnector/config/validate \ + -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.test 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..e0e6ff4 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -174,10 +174,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 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 From d0301d137b1768218cb71dfa6075899e1a4f4038 Mon Sep 17 00:00:00 2001 From: Shadow Date: Sat, 9 Aug 2025 21:19:08 +0700 Subject: [PATCH 2/9] fix: yamllint and trivy scan error --- .github/workflows/kafka-ci.yml | 5 +++++ docker-compose.yml | 20 ++++++++++++++------ 2 files changed, 19 insertions(+), 6 deletions(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index 5ef5487..cba60f8 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -1,5 +1,10 @@ name: Kafka CDC Pipeline CI +permissions: + actions: read + contents: read + security-events: write + on: push: branches: diff --git a/docker-compose.yml b/docker-compose.yml index e0e6ff4..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 @@ -225,4 +233,4 @@ volumes: networks: lakepulse_network: - driver: bridge \ No newline at end of file + driver: bridge From e8b63b9f2596f62b0f14538c0d7907cf3c455b5f Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 09:53:54 +0700 Subject: [PATCH 3/9] fix: format lint for kafka-ci.yml --- .github/workflows/kafka-ci.yml | 34 +++++++++++++++++++++++----------- 1 file changed, 23 insertions(+), 11 deletions(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index cba60f8..f31b963 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -1,3 +1,4 @@ +--- name: Kafka CDC Pipeline CI permissions: @@ -87,7 +88,7 @@ jobs: permissions: contents: read packages: write - + outputs: image-digest: ${{ steps.build.outputs.digest }} image-tag: ${{ steps.meta.outputs.tags }} @@ -173,11 +174,13 @@ jobs: - name: Start test services run: | # Update docker-compose to use test image - sed -i 's|lakepulse/kafka-connect:latest|${{ needs.build.outputs.image-tag }}|' docker-compose.yml - + IMAGE_TAG="${{ needs.build.outputs.image-tag }}" + sed -i "s|lakepulse/kafka-connect:latest|${IMAGE_TAG}|" \ + docker-compose.yml + # Start minimal services for testing docker compose --env-file .env.test up -d broker schema-registry minio - + # Wait for services to be ready sleep 30 @@ -185,25 +188,34 @@ jobs: run: | # Start Kafka Connect with test image docker compose --env-file .env.test up -d kafka-connect - + # Wait for Connect to start - timeout 120 bash -c 'until curl -f http://localhost:8083/; do sleep 5; done' + 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"))' + 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 - curl -X PUT http://localhost:8083/connector-plugins/io.debezium.connector.postgresql.PostgresConnector/config/validate \ + 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 - curl -X PUT http://localhost:8083/connector-plugins/io.confluent.connect.s3.S3SinkConnector/config/validate \ + 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' From 0ff8b17390b387b80dd1acef59bd4ec7d4662a5d Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 09:59:28 +0700 Subject: [PATCH 4/9] fix: missing .env in actions --- .env.example | 9 ++++++++- .github/workflows/kafka-ci.yml | 3 +++ 2 files changed, 11 insertions(+), 1 deletion(-) 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 index f31b963..d65a49b 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -37,6 +37,9 @@ jobs: - 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: From 7a3cf32a025b2b89b38939432f50aff436d81a12 Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 10:16:38 +0700 Subject: [PATCH 5/9] fix: missing env file again --- .github/workflows/kafka-ci.yml | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index d65a49b..9015646 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -164,15 +164,7 @@ jobs: - name: Create test environment file run: | - cat > .env.test << EOF - POSTGRES_USER=${{ secrets.POSTGRES_USER }} - POSTGRES_PASSWORD=${{ secrets.POSTGRES_PASSWORD }} - POSTGRES_DB=${{ secrets.POSTGRES_DB }} - POSTGRES_PORT=${{ secrets.POSTGRES_PORT }} - MINIO_ROOT_USER=${{ secrets.MINIO_ROOT_USER }} - MINIO_ROOT_PASSWORD=${{ secrets.MINIO_ROOT_PASSWORD }} - KAFKA_BOOTSTRAP_SERVERS=${{ secrets.KAFKA_BOOTSTRAP_SERVERS }} - EOF + copy .env.example .env - name: Start test services run: | @@ -182,7 +174,8 @@ jobs: docker-compose.yml # Start minimal services for testing - docker compose --env-file .env.test up -d broker schema-registry minio + docker compose --env-file .env up -d postgres broker \ + schema-registry minio # Wait for services to be ready sleep 30 @@ -190,7 +183,7 @@ jobs: - name: Test Kafka Connect image run: | # Start Kafka Connect with test image - docker compose --env-file .env.test up -d kafka-connect + docker compose --env-file .env up -d kafka-connect # Wait for Connect to start timeout 120 bash -c \ @@ -225,5 +218,5 @@ jobs: - name: Cleanup test environment if: always() run: | - docker compose --env-file .env.test down -v + docker compose --env-file .env down -v docker system prune -f From 972905aee8e3a90be1b80ad9c3de8aa186b11ac4 Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 10:20:36 +0700 Subject: [PATCH 6/9] fix: remove unnecessary postgres image --- .github/workflows/kafka-ci.yml | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index 9015646..0f67eb9 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -143,20 +143,6 @@ jobs: name: Integration Tests runs-on: ubuntu-latest needs: [build] - services: - postgres: - image: postgres:15-alpine - env: - POSTGRES_USER: ${{ secrets.POSTGRES_USER }} - POSTGRES_PASSWORD: ${{ secrets.POSTGRES_PASSWORD }} - POSTGRES_DB: ${{ secrets.POSTGRES_DB }} - options: >- - --health-cmd pg_isready - --health-interval 10s - --health-timeout 5s - --health-retries 5 - ports: - - 5432:5432 steps: - name: Checkout repository From 666a46fcab82eae42d80efb3504f0bea248d69fd Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 10:22:26 +0700 Subject: [PATCH 7/9] fix: wrong copy command XD --- .github/workflows/kafka-ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index 0f67eb9..7373cb0 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -150,7 +150,7 @@ jobs: - name: Create test environment file run: | - copy .env.example .env + cp .env.example .env - name: Start test services run: | From 2f22eacf6888d569f977610dce6400f7a3bdd020 Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 10:29:24 +0700 Subject: [PATCH 8/9] fix: multiple image-tag --- .github/workflows/kafka-ci.yml | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index 7373cb0..ce95c96 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -154,10 +154,9 @@ jobs: - name: Start test services run: | - # Update docker-compose to use test image - IMAGE_TAG="${{ needs.build.outputs.image-tag }}" - sed -i "s|lakepulse/kafka-connect:latest|${IMAGE_TAG}|" \ - docker-compose.yml + # 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 \ From f1958acc7b1a892e4fb9004eeddce58593082635 Mon Sep 17 00:00:00 2001 From: Shadow Date: Mon, 11 Aug 2025 10:31:10 +0700 Subject: [PATCH 9/9] fix: yamllint error again... --- .github/workflows/kafka-ci.yml | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/.github/workflows/kafka-ci.yml b/.github/workflows/kafka-ci.yml index ce95c96..3f8cc04 100644 --- a/.github/workflows/kafka-ci.yml +++ b/.github/workflows/kafka-ci.yml @@ -156,7 +156,8 @@ jobs: 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 + 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 \