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.
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_ONCEmode (duplicate-delivery risk); missing checkpoint timeout. - FlinkStateAnalyzer — flags heap-bound state backends
(
HashMapStateBackend/MemoryStateBackend) that risk OOM at scale, RocksDB without incremental checkpoints, deprecatedFsStateBackend, andkey_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
WatermarkStrategyassigned (windows may never fire), excessive bounded-out-of-orderness, and keyed streams with no windowing.
- FlinkCheckpointAnalyzer — flags stateful pipelines with no
- Apache Spark — real, first-class analyzers (
--framework spark), same approach, run against a structured IR extracted from PySpark source:- SparkShuffleAnalyzer — flags
spark.sql.shuffle.partitionsleft 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
autoBroadcastJoinThresholddisabled (-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
writeStreamqueries with nocheckpointLocation(no recovery guarantee), no explicit trigger,cache()/persist()with no matchingunpersist(), and the classic Structured Streaming pitfall of agroupByaggregation inappendoutput mode with no watermark (fails at query start in real Spark).
- SparkShuffleAnalyzer — flags
- 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
- Build Summary - Phase implementation details and history
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.ymlworkflow'scargo build --workspace --verbosestep tried to link the PyO3extension-modulebindings crate (bindings/python) as a plain cdylib outside ofmaturin--ld: symbol(s) not found for architecture arm64on the macOS runner (reproduced locally). Because that step ran beforeRun tests, the test step never executed in CI. Fix: scope both theBuild workspaceandRun testssteps to-p pybeamguard-coreinstead of--workspace, since that crate is the only one meant to be built by plaincargo build/cargo test--bindings/pythonis only ever built viamaturin, which the separatepython-bindingsCI job already exercises end-to-end. Verified locally with the exact new CI commands:cargo build -p pybeamguard-coreandcargo test -p pybeamguard-coreboth pass (73 unit + 11 integration tests), as doescargo fmt --all -- --checkandcargo 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 viamaturin, not this workflow, andpip install pybeamguard+pybeamguard analyzewere 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/15MBwith 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 viagrep, not a fullcargo test/pytestrun — see the CI badge for the authoritative, current count). - No open GitHub issues and no
TODO/FIXME/XXXmarkers found incrates/orsrc/as of this pass. analyzeonly 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
- GitHub Issues: https://github.com/Mullassery/PyBeamGuard/issues
- Repository: https://github.com/Mullassery/PyBeamGuard
- PyPI: https://pypi.org/project/pybeamguard/
- Email: mullassery@gmail.com
- Author: @Mullassery
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:
- Manual review (slow, inconsistent)
- Running expensive test jobs (costly, time-consuming)
- 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)
| File | Size | Uploaded | |
|---|---|---|---|
| pybeamguard-1.2.0.tar.gz | 91.3 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|