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))