Skip to main content

PyBeamGuard

Catch Apache Beam failures before deployment. Forecast costs. Fix hot keys.

Analyze Beam pipelines pre-deployment to identify bottlenecks, reliability risks, and cost drivers. FREE. No cloud account required. Works offline.

A static analysis platform for Apache Beam, Apache Flink, and Apache Spark pipelines.

CI Version License PyPI


Comparison with Similar Tools

Why PyBeamGuard?

Feature PyBeamGuard Vendor Profiler Managed Runner Console
Pre-deployment analysis Yes No No
Cost forecasting Heuristic estimate (e.g. $48-$2.5K/mo range) No Post-deploy only
Hot key detection Yes No No
Shuffle analysis Yes No Post-deploy only
Windowing validation Yes No No
Cost FREE Included with cloud subscription Included with cloud subscription
Setup required None Cloud account Cloud account
Offline capable Yes No No

Bottom line: Pre-deployment analysis you control, costs you forecast before running, no vendor lock-in.


Quick Start

Installation

Requires Python 3.10 or later

# Using pip
pip install pybeamguard

# Using uv (faster)
uv pip install pybeamguard

# Verify installation
pybeamguard --version

Analyze a Pipeline

# Text output (default)
pybeamguard analyze pipeline.py

# JSON output  
pybeamguard analyze pipeline.py --format json

# With data profile
pybeamguard analyze pipeline.py --data-profile profile.json

Example Output

=== PyBeamGuard Analysis Report ===

Overall Risk Score: 78/100
Total Findings: 5

CRITICAL ISSUES
• Hot key probability detected on customer_id aggregation

HIGH PRIORITY ISSUES
• Large shuffle stage in join operation
• Missing dead-letter queue on parse failures

MEDIUM PRIORITY ISSUES
• Unbounded state growth risk

Estimated Cost: $2,300/month → Optimized: $1,350/month (41% savings)

Features

All Features FREE - Proprietary Software

10 Intelligent Analyzers

Analyzer Purpose Version
Graph Intelligence Extract pipeline topology, detect cycles v0.1
Hot Key Detection Identify key skew & worker imbalance v0.1
Shuffle Analysis Quantify expensive shuffle operations v0.1
Windowing Validation Ensure streaming correctness v0.1
State Auditor Prevent state-related failures v0.1
Cost Intelligence Forecast managed-runner spend v0.1
Reliability Analysis Detect operational weaknesses v0.1
Best Practices Engine Rule-based Beam optimization checks v0.1
Deployment Auditor Worker sizing & config validation v0.1
Architecture Review Executive summary & synthesis v0.1

Framework Support (Free)

  • Apache Beam — full pipeline graph extraction + all 10 analyzers
  • Apache Flink — real, first-class analyzers (--framework flink), same regex/line-based static-analysis approach as Beam, run against a structured IR extracted from PyFlink source:
    • FlinkCheckpointAnalyzer — flags stateful pipelines with no enable_checkpointing(...) at all; checkpoint intervals that are too aggressive (<1s, barrier-alignment overhead) or too long (>10min, large replay window on restart); AT_LEAST_ONCE mode (duplicate-delivery risk); missing checkpoint timeout.
    • FlinkStateAnalyzer — flags heap-bound state backends (HashMapStateBackend/MemoryStateBackend) that risk OOM at scale, RocksDB without incremental checkpoints, deprecated FsStateBackend, and key_by(...) expressions on high-skew-risk domains (customer/ tenant/user/... — the Flink analog of Beam's hot-key detection, applied to keyed state).
    • FlinkWatermarkAnalyzer — flags event-time windows with no WatermarkStrategy assigned (windows may never fire), excessive bounded-out-of-orderness, and keyed streams with no windowing.
  • Apache Spark — real, first-class analyzers (--framework spark), same approach, run against a structured IR extracted from PySpark source:
    • SparkShuffleAnalyzer — flags spark.sql.shuffle.partitions left at the 200 default alongside multiple wide transforms, set too low (OOM/ skew risk) or too high (per-task overhead), and jobs with several shuffle-triggering transforms (groupBy/join/distinct/ repartition/coalesce/orderBy).
    • SparkJoinAnalyzer — flags autoBroadcastJoinThreshold disabled (-1) or set dangerously high (broadcast OOM risk), and join keys on high-skew-risk domains (the Spark analog of Beam's hot-key detection, applied to shuffle join keys).
    • SparkStreamingAnalyzer — flags writeStream queries with no checkpointLocation (no recovery guarantee), no explicit trigger, cache()/persist() with no matching unpersist(), and the classic Structured Streaming pitfall of a groupBy aggregation in append output mode with no watermark (fails at query start in real Spark).
  • Kafka Streams (not started)
  • Ray Data (not started)

Both are driven by the same analyze command:

pybeamguard analyze streaming_job.py --framework flink
pybeamguard analyze etl.py --framework spark

Configurable Rule Engine (YAML)

The Hot Key, Cost, Spark Join, and State analyzers' rule data (risk-pattern keyword lists, cardinality/size thresholds, cost rates) is no longer compiled-in const data — it's loaded from a RulesConfig, which defaults to the exact values that used to be hardcoded but can be overridden with a YAML file:

# rules.yaml — every section/field is optional; anything omitted keeps its
# built-in default (see PARTIAL override semantics below).
hotkey:
  high_risk_patterns: ["customer", "tenant", "region"]  # add your own domain keywords
  low_cardinality_threshold: 500
cost:
  worker_machine_cost_per_hour: 0.42  # match your actual machine type's rate
spark_join:
  broadcast_threshold_high_risk_bytes: 536870912
state:
  large_measured_state_size_gb: 250.0
pybeamguard analyze pipeline.py --rules rules.yaml
pybeamguard analyze etl.py --framework spark --rules rules.yaml

A file only needs to set the fields it wants to change — untouched sections/fields keep their default. This is available both from the Rust CLI (--rules <path>) and the Python bindings (analyze(code, rules_yaml=...), analyze_spark(...), get_json_report(...), etc. all accept an optional rules_yaml string). Flink's analyzer suite (checkpointing/state-backend/watermark) has no rule-driven thresholds, so --rules has no effect with --framework flink.


Use Cases

For Data Engineers

Pre-deployment validation: "Will this scale? What will it cost?"

pybeamguard analyze my_pipeline.py
# Identifies 3 hot key risks
# Estimates $850/month cost
# Warns of unbounded state growth

For Platform Teams

CI/CD enforcement: fail the build if any finding is critical. analyze takes a single pipeline file (there's no built-in directory/glob support yet), so scanning a directory of pipelines means looping over the files:

# .github/workflows/pipeline-validation.yml
- run: |
    for f in pipelines/*.py; do
      pybeamguard analyze "$f" --fail-on critical || exit 1
    done

For FinOps Teams

Cost attribution: "Why is this pipeline $2,500/month?"

pybeamguard analyze pipeline.py --format json | jq '.[] | select(.analyzer_name=="CostAnalyzer")'
# "estimated_total_cost_per_month": 2500.00
# "estimated_shuffle_cost_per_month": 1500.00  ← Cost hotspot

Note: these dollar figures are a rough, heuristic pre-deployment estimate, not a validated billing forecast — see the Cost Forecasting caveat below.


Installation

From PyPI (Recommended)

Python 3.10+ with pip or uv:

# Using pip
pip install pybeamguard

# Using uv
uv pip install pybeamguard

# Verify installation
pybeamguard --version

From GitHub Releases (standalone Rust binary, no Python required)

GitHub Releases publishes the standalone pybeamguard Rust CLI binary for macOS, Linux, and Windows — this is a separate build from the PyPI wheel above and has no Python dependency at all:

# Download the binary for your platform from:
# https://github.com/Mullassery/PyBeamGuard/releases
chmod +x pybeamguard-macos-universal   # or the linux/windows artifact
./pybeamguard-macos-universal --version

From Source

Requires Rust 1.70+ and Python 3.10+:

git clone https://github.com/Mullassery/PyBeamGuard.git
cd PyBeamGuard

# Python package (PyO3 bindings), installed into your active venv:
pip install maturin
maturin develop
pybeamguard --version

# OR the standalone Rust CLI binary (no Python involved):
cargo build --release --bin pybeamguard
./target/release/pybeamguard --version

Documentation


Examples

Example 1: Simple Batch Pipeline

# pipeline.py
import apache_beam as beam

with beam.Pipeline() as p:
    result = (
        p
        | 'Read' >> beam.io.ReadFromText('input.txt')
        | 'Parse' >> beam.ParDo(ParseFn())
        | 'GroupByCustomer' >> beam.GroupByKey()
        | 'CountPerCustomer' >> beam.CombinePerKey(sum)
        | 'Write' >> beam.io.WriteToText('output.txt')
    )
$ pybeamguard analyze pipeline.py

HIGH PRIORITY
• High hot-key probability on customer_id
  Impact: 3-5x latency increase
  Mitigation: Apply key sharding strategy

Cost Estimate
  Compute: $18/month
  Shuffle: $30/month
  Total: $48/month

Recommendation: Implement key sharding before production

Example 2: With Data Profile

// profile.json
{
  "estimated_throughput_per_sec": 10000,
  "average_element_size_bytes": 500,
  "key_cardinality": 50000,
  "estimated_state_size_gb": 5.0
}
$ pybeamguard analyze pipeline.py --data-profile profile.json --format json

{
  "analyzer_name": "CostAnalyzer",
  "findings": [...],
  "metrics": {
    "estimated_total_cost_per_month": 2350.00,
    "estimated_compute_cost_per_month": 175.00,
    "estimated_shuffle_cost_per_month": 900.00,
    "estimated_state_cost_per_month": 1275.00
  }
}

Performance

No committed benchmark script or results file backs a specific latency/memory number, so none is stated here as fact — see the Release Status and FAQ sections below for the honest version ("not independently benchmarked at scale"). What's verifiable today:

Metric Value
Tests ~78 Rust unit tests + 13 integration tests + 23 Python binding/CLI tests (see CI badge for the exact, current count)

Requirements

  • macOS 10.13+ (Intel/Apple Silicon)
  • Linux (glibc 2.31+, x86_64)
  • Windows 10/11 (x86_64)

The standalone Rust binary (GitHub Releases) has no Python runtime, dependencies, or environment variables required. The PyPI package (pip install pybeamguard) requires Python 3.10+, same as any Python package — it's a compiled PyO3 extension module with a thin Python shim, not a pure-Python implementation.


Contributing

Contributions welcome!

# Build the Rust core + CLI binary
cargo build --release --bin pybeamguard
# Note: plain `cargo build --workspace` (building the PyO3 bindings crate
# without maturin) still fails to link on macOS -- that crate's
# `extension-module` feature needs maturin's dynamic-lookup linker flags,
# which a bare `cargo build` doesn't supply. CI works around this by
# scoping its build/test steps to `-p pybeamguard-core` (see
# .github/workflows/ci.yml); do the same locally, or use
# `maturin develop`/`maturin build` for the Python bindings instead
# (below), which is what the real PyPI release uses.

# Run the Rust test suite (unit + integration)
cargo test -p pybeamguard-core

# Lint
cargo fmt --all -- --check
cargo clippy --workspace -- -D warnings

# Build + install the Python bindings locally, then run the Python test suite
pip install maturin pytest
maturin develop
pytest tests/

# Analyze an example pipeline
./target/release/pybeamguard analyze examples/pipeline_simple.py

Release Status

Current state: working proof-of-concept for Apache Beam, Flink, and Spark analysis, with packaging and CI around it. Concretely, what's implemented and tested today:

  • 10 intelligent analyzers over Apache Beam pipelines (regex/heuristic-based, not full AST analysis)
  • 3 intelligent analyzers over Apache Flink pipelines (checkpointing, state backend, watermark/windowing) and 3 over Apache Spark pipelines (shuffle partitioning, join/broadcast strategy, streaming checkpoint/trigger/output-mode) — see Framework Support above
  • Python bindings via PyO3 abi3, real pip install-able package
  • --fail-on <severity> CI gating and --data-profile-informed cost/hot-key estimates
  • Rust unit + integration tests, Python binding/CLI tests -- verified by running them directly (cargo test -p pybeamguard-core: 78 unit + 13 integration tests passing; pytest tests/: 23 Python tests passing), and CI is green running the same commands
  • <500ms analysis per pipeline (small/medium pipelines; not independently benchmarked at scale)

Explicitly not implemented (removed from this codebase to stop overclaiming rather than left as unused/untested scaffolding): organization governance (cost budgets, SLOs, policy enforcement), audit logging, and Airflow/dbt/data-contract/FinOps ecosystem integrations. These were previously present as struct definitions with no wiring into the actual analysis path and no way to test them without external systems this project doesn't have access to.

Future Roadmap (aspirational, not started):

  • Kafka Streams, Ray Data framework support
  • Directory/glob input to analyze (currently single-file only)
  • Re-introduce org governance / audit logging as real, tested features if there's demand
  • Python plugin hooks / Rego/OPA integration for rule logic (not just rule data) — out of scope for now. Done: rule data externalization (HIGH_RISK_PATTERNS/MEDIUM_RISK_PATTERNS/cost rates/thresholds) via YAML — see "Configurable Rule Engine (YAML)" above.

Known Issues

  • CI was red from 2026-08-07 to 2026-08-23 (fixed in this pass). Root cause: the ci.yml workflow's cargo build --workspace --verbose step tried to link the PyO3 extension-module bindings crate (bindings/python) as a plain cdylib outside of maturin -- ld: symbol(s) not found for architecture arm64 on the macOS runner (reproduced locally). Because that step ran before Run tests, the test step never executed in CI. Fix: scope both the Build workspace and Run tests steps to -p pybeamguard-core instead of --workspace, since that crate is the only one meant to be built by plain cargo build/cargo test -- bindings/python is only ever built via maturin, which the separate python-bindings CI job already exercises end-to-end. Verified locally with the exact new CI commands: cargo build -p pybeamguard-core and cargo test -p pybeamguard-core both pass (73 unit + 11 integration tests), as does cargo fmt --all -- --check and cargo clippy --workspace -- -D warnings (clippy doesn't need the final cdylib link, so it's safe to leave workspace-wide). The PyPI package itself was never affected -- it's built via maturin, not this workflow, and pip install pybeamguard + pybeamguard analyze were confirmed working end-to-end.
  • No committed benchmark script or results file backs a latency/memory/binary-size number, so the Performance section above no longer states one as fact — a previous version of this README claimed <500ms / <50MB / 15MB with nothing checked in to reproduce those figures.
  • The Rust/Python test counts previously stated in this README (29 Rust unit tests + 7 integration tests) were stale; the current source has roughly 78 #[test]-annotated Rust unit tests, 13 Rust integration tests, and 23 Python tests (counted via grep, not a full cargo test/pytest run — see the CI badge for the authoritative, current count).
  • No open GitHub issues and no TODO/FIXME/XXX markers found in crates/ or src/ as of this pass.
  • analyze only accepts a single pipeline file; directory/glob scanning requires the loop shown in the Platform Teams example above (tracked in Future Roadmap).

License

Proprietary Software — FREE forever, no licensing tiers, no paywalls.

See LICENSE file for complete terms. All features available to all users.

Use Cases:

  • Commercial use
  • Internal tools
  • Research
  • Education
  • Open source projects

Support & Contact


FAQ

Q: How much does PyBeamGuard cost?
A: FREE. PyBeamGuard is proprietary software with no licensing fees, no tiers, no paywalls. All features available to everyone.

Q: Does PyBeamGuard require Python?
A: Depends how you install it. The standalone Rust binary from GitHub Releases has zero dependencies — download and run. The pip install pybeamguard package is a compiled PyO3 extension with a thin Python shim, so it requires Python 3.10+ like any Python package.

Q: What pipeline sizes can it analyze?
A: The parser and analyzers are simple regex/line-based passes over the source, so there's no architectural node-count ceiling, but this hasn't been independently benchmarked at large scale (e.g. 1,000+ node pipelines).

Q: How accurate are the cost estimates?
A: Treat them as an order-of-magnitude planning signal, not a validated billing forecast — the underlying cost model uses simplified, only partially-verified pricing assumptions (see CostAnalyzer's doc comments in crates/core/src/analyzers/cost.rs). Supplying --data-profile replaces some flat per-operation guesses with your real throughput/state figures, which narrows the estimate, but doesn't make it a guarantee. Always confirm with your cloud provider's pricing calculator or a real test run before committing to a budget.

Q: Can I use this in CI/CD?
A: Yes — pybeamguard analyze pipeline.py --fail-on critical exits non-zero if any finding meets or exceeds the given severity, so it works in any CI system with a shell (GitHub Actions, GitLab CI, Jenkins, Cloud Build, ...). There's no bundled CI-specific plugin/action, just a CLI with a meaningful exit code.

Q: What about Spark, Flink, Kafka Streams?
A: Spark and Flink both have real, dedicated analyzer suites today — run with pybeamguard analyze <file> --framework flink or --framework spark (see Framework Support above for exactly what each analyzer checks). Kafka Streams and Ray Data support hasn't been started.


Detailed Comparison Matrix

Analysis Capabilities

Capability PyBeamGuard Beam Native Tools Managed Runner Console Monitoring Tools
Pipeline graph extraction Yes No No No
Complexity scoring Yes No No No
Hot key detection Yes Heuristic (keyword pattern + optional measured cardinality) No (disabled 2022) Partial Disabled for streaming No
Shuffle quantification Yes Per-stage (rule-based) No Partial Aggregate only Partial Post-deploy only
State growth prediction Yes Heuristic (flags stateful ops) No No No
Cost forecasting Partial Pre-deploy, rough heuristic No Partial Post-deploy estimate No
Best practices engine Yes Rule-based checks No No No
Deployment audit Yes No No No
Architecture review Yes Rule-based weighted scoring (not ML) No No Partial Manual only

Deployment & Integration

Aspect PyBeamGuard Vendor Profiler Managed Runner Console
Installation pip install / wheel Built-in (cloud vendor) Built-in (cloud vendor)
Setup time <1 minute Account required Account required
Offline support Yes Full No No No No
CI/CD gating Yes --fail-on <severity> exit code (works with any CI system) No No No No
Python version 3.10+ (via PyO3) Any (cloud vendor) Any (cloud vendor)
Platform support macOS, Linux, Windows Cloud vendor only Cloud vendor only

Cost

Feature PyBeamGuard Competitors
Tool cost FREE Managed runner console: free (but runs expensive test jobs)
Cost forecasting Partial Rough, pre-deploy heuristic (see caveat below) No Requires running pipelines
Test job cost Yes Save money by not needing a real run for a first pass No Must run to estimate cost

Time to Insight

Task PyBeamGuard Vendor Profiler Managed Runner Console
Analyze pipeline <1 sec N/A (need to run) N/A (need to run)
Detect hot keys <1 sec 30+ min (with run) 30+ min (with run)
Forecast cost <1 sec N/A 24-48 hours (post-deploy)
Architecture review <2 sec N/A N/A

Why PyBeamGuard Exists

PyBeamGuard fills a critical gap:

The Problem: The major managed Beam runner disabled hot key detection for streaming pipelines in March 2022. No other tool provides pre-deployment Beam analysis. Teams are left with:

  1. Manual review (slow, inconsistent)
  2. Running expensive test jobs (costly, time-consuming)
  3. Production incidents (expensive, damaging)

The Solution: PyBeamGuard brings expert-level Beam analysis to every team, offline and for free.


Built for data engineers everywhere.

Release files for pybeamguard 1.2.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for pybeamguard 1.2.0
File Size Uploaded
pybeamguard-1.2.0.tar.gz 91.3 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for pybeamguard 1.2.0
File Interpreter ABI Platform
pybeamguard-1.2.0-cp311-abi3-macosx_11_0_arm64.whl CPython 3.11 abi3 macOS 11.0+ ARM64 Details

Total release size: 1.3 MB

Release files / pybeamguard-1.2.0.tar.gz

Download URL pybeamguard-1.2.0.tar.gz
Size 91.3 kB
Tags Source
SHA-256 checksum
How to use checksums
e300fa5b6a96e8bb4a461d26cecf3a8f317b85a0249841e05c96774e852f9242
BLAKE2b-256 checksum
How to use checksums
4c91bbeb673d55c8508b1da8760c1ff4fedccafcad84aa24f705ddd79c432112
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.16

Release files / pybeamguard-1.2.0-cp311-abi3-macosx_11_0_arm64.whl

Download URL pybeamguard-1.2.0-cp311-abi3-macosx_11_0_arm64.whl
Size 1.3 MB
Tags CPython 3.11 abi3 macOS 11.0+ ARM64
SHA-256 checksum
How to use checksums
0560e92f852f1e132ed6e689fc8fedbbec0d8c6a0527e82f87364efdae7458ed
BLAKE2b-256 checksum
How to use checksums
bec5e46d23b9ba65432ac7e9f65c866209196d8468e2cacee9c17897f15edeb5
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.16

Release history Release notifications | RSS feed

1.2.1

2 release files

This release

1.2.0 This release

2 release files

1.1.2

1 release file

1.1.1

1 release file

1.1.0

1 release file

1.0.0

1 release file

0.9.0

1 release file

0.8.0

1 release file

0.7.0

1 release file

0.6.0

1 release file

0.5.0

1 release file

0.4.0

1 release file

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page