Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
KAFKA_SCHEMA_REGISTRY_URL=http://schema-registry:8081
CDC_TOPIC_PREFIX=cdc
208 changes: 208 additions & 0 deletions .github/workflows/kafka-ci.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,208 @@
---
name: Kafka CDC Pipeline CI

permissions:
actions: read
contents: read
security-events: write

on:

Check warning on line 9 in .github/workflows/kafka-ci.yml

View workflow job for this annotation

GitHub Actions / Lint Configuration Files

9:1 [truthy] truthy value should be one of [false, true]

Check warning on line 9 in .github/workflows/kafka-ci.yml

View workflow job for this annotation

GitHub Actions / Lint Configuration Files

9:1 [truthy] truthy value should be one of [false, true]
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
128 changes: 126 additions & 2 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -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 \
Expand All @@ -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)"
@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)"
Loading
Loading