CI/CD для Spark: тестирование, линтинг и автоматический деплой

Production Spark pipeline должен проходить автоматическую проверку: unit-тесты, type checking, Mermaid-валидация, сборка артефакта и деплой через GitHub Actions.

platform

Зачем CI/CD для Spark

Без автоматизации:

  • Баги находятся только в production (после часов выполнения)
  • «Работает на моей машине» → падает на кластере (разные версии Python/PySpark)
  • Нет gate перед мержем → сломанный код попадает в main

Структура Spark проекта для CI/CD

spark-project/
├── jobs/
│   ├── etl_pipeline.py        # основной job
│   └── transformations.py     # чистые функции (тестируемые)
├── tests/
│   ├── conftest.py            # SparkSession fixture
│   ├── test_transformations.py
│   └── test_integration.py
├── requirements.txt
├── requirements-dev.txt       # pytest, chispa, mypy, ruff
├── Makefile
└── .github/workflows/
    ├── test.yml               # CI: lint + test
    └── deploy.yml             # CD: package + submit

GitHub Actions: CI pipeline

# .github/workflows/test.yml
name: Spark CI

on:
  pull_request:
    branches: [main]
  push:
    branches: [main]

jobs:
  lint-and-test:
    runs-on: ubuntu-latest

    steps:

      - uses: actions/checkout@v4

      - uses: actions/setup-python@v5
        with:
          python-version: "3.11"

      - uses: actions/setup-java@v4
        with:
          java-version: "17"
          distribution: "temurin"

      - name: Install dependencies
        run: |
          pip install -r requirements-dev.txt
          pip install pyspark==3.5.0

      - name: Lint (ruff)
        run: ruff check jobs/ tests/

      - name: Type check (mypy)
        run: mypy jobs/ --ignore-missing-imports

      - name: Unit tests
        run: |
          pytest tests/ -v \
            --tb=short \
            --junit-xml=test-results.xml
        env:
          PYSPARK_PYTHON: python3
          JAVA_HOME: /usr/lib/jvm/temurin-17-amd64

      - name: Upload test results
        uses: actions/upload-artifact@v4
        if: always()
        with:
          name: test-results
          path: test-results.xml

GitHub Actions: CD pipeline (деплой на YARN/K8s)

# .github/workflows/deploy.yml
name: Spark Deploy

on:
  push:
    branches: [main]
    paths:

      - 'jobs/**'

jobs:
  package-and-deploy:
    runs-on: ubuntu-latest
    environment: production

    steps:

      - uses: actions/checkout@v4

      - uses: actions/setup-python@v5
        with:
          python-version: "3.11"

      - name: Package Python dependencies
        run: |
          pip install venv-pack
          python -m venv .venv
          source .venv/bin/activate
          pip install -r requirements.txt
          venv-pack -o spark-env.tar.gz

      - name: Upload artifacts to S3
        run: |
          aws s3 cp spark-env.tar.gz s3://$BUCKET/artifacts/spark-env-${{ github.sha }}.tar.gz
          aws s3 cp jobs/etl_pipeline.py s3://$BUCKET/jobs/etl_pipeline-${{ github.sha }}.py
        env:
          BUCKET: my-spark-artifacts
          AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
          AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}

      - name: Submit Spark job (YARN cluster mode)
        run: |
          spark-submit \
            --master yarn \
            --deploy-mode cluster \
            --archives s3://$BUCKET/artifacts/spark-env-${{ github.sha }}.tar.gz#env \
            --conf spark.pyspark.python=./env/bin/python \
            --conf spark.app.version=${{ github.sha }} \
            s3://$BUCKET/jobs/etl_pipeline-${{ github.sha }}.py
        env:
          BUCKET: my-spark-artifacts

Makefile для локальной разработки

# Makefile
.PHONY: test lint typecheck clean

test:
    pytest tests/ -v --tb=short

lint:
    ruff check jobs/ tests/
    ruff format --check jobs/ tests/

typecheck:
    mypy jobs/ --ignore-missing-imports

format:
    ruff format jobs/ tests/

package:
    venv-pack -o spark-env.tar.gz

clean:
    find . -type d -name __pycache__ -exec rm -rf {} +
    rm -f spark-env.tar.gz test-results.xml

ci: lint typecheck test

Versioning артефактов

# jobs/etl_pipeline.py - версия передаётся через конфиг
import os
from pyspark.sql import SparkSession

spark = SparkSession.getOrCreate()

app_version = spark.conf.get("spark.app.version", "unknown")
spark.sparkContext.setJobDescription(f"ETL Pipeline v{app_version}")

# Логировать версию в метриках
spark._jvm.org.apache.log4j.LogManager \
    .getLogger("ETL") \
    .info(f"Starting ETL Pipeline version: {app_version}")

Идемпотентность деплоя

# ✅ Каждый деплой должен быть безопасен для повторного запуска
# Использовать OVERWRITE вместо APPEND для partitioned таблиц:
df.write \
    .mode("overwrite") \
    .option("partitionOverwriteMode", "dynamic") \
    .parquet("s3://output/events/")

# ✅ Checkpoint при повторном запуске:
CHECKPOINT_PATH = f"s3://checkpoints/{app_version}/"

Quality gate в CI

# tests/test_data_quality.py - запускается в CI против sample данных
def test_no_nulls_in_pk(spark):
    df = transform(create_test_df(spark))
    null_count = df.filter(col("user_id").isNull()).count()
    assert null_count == 0, f"Found {null_count} null user_ids"

def test_row_count_not_empty(spark):
    df = transform(create_test_df(spark))
    assert df.count() > 0

def test_schema_matches_contract(spark):
    df = transform(create_test_df(spark))
    expected_cols = {"user_id", "email", "amount", "event_date"}
    assert expected_cols.issubset(set(df.columns))