Metadata-Version: 2.4
Name: pybeamguard
Version: 1.2.0
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: Intended Audience :: System Administrators
Classifier: Natural Language :: English
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Programming Language :: Rust
Classifier: Topic :: Software Development :: Libraries
Classifier: Topic :: System :: Monitoring
License-File: LICENSE
Summary: Apache Beam & Dataflow pipeline analysis, reliability validation, and cost optimization platform
Keywords: apache-beam,dataflow,pipeline,analysis,optimization,cost
Author-email: Georgi Mammen Mullassery <mullassery@gmail.com>
License: Proprietary License — Free to use with explicit attribution
Requires-Python: >=3.10
Description-Content-Type: text/markdown; charset=UTF-8; variant=GFM
Project-URL: Documentation, https://github.com/Mullassery/pybeamguard#readme
Project-URL: Homepage, https://github.com/Mullassery/pybeamguard
Project-URL: Issues, https://github.com/Mullassery/pybeamguard/issues
Project-URL: Repository, https://github.com/Mullassery/pybeamguard.git

# 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](https://github.com/Mullassery/PyBeamGuard/actions/workflows/ci.yml/badge.svg)](https://github.com/Mullassery/PyBeamGuard/actions/workflows/ci.yml)
[![Version](https://img.shields.io/badge/version-1.1.2-blue)](https://github.com/Mullassery/PyBeamGuard/releases)
[![License](https://img.shields.io/badge/license-Proprietary-blue)](LICENSE)
[![PyPI](https://img.shields.io/badge/PyPI-pybeamguard-blue)](https://pypi.org/project/pybeamguard/)

---

## 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**

```bash
# Using pip
pip install pybeamguard

# Using uv (faster)
uv pip install pybeamguard

# Verify installation
pybeamguard --version
```

### Analyze a Pipeline

```bash
# 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:

```bash
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:

```yaml
# 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
```

```bash
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?"

```bash
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:

```yaml
# .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?"

```bash
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:

```bash
# 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:

```bash
# 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+:

```bash
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](BUILD_SUMMARY.md)** - Phase implementation details and history

---

## Examples

### Example 1: Simple Batch Pipeline

```python
# 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')
    )
```

```bash
$ 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

```json
// profile.json
{
  "estimated_throughput_per_sec": 10000,
  "average_element_size_bytes": 500,
  "key_cardinality": 50000,
  "estimated_state_size_gb": 5.0
}
```

```bash
$ 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!

```bash
# 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](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](https://github.com/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:
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.**

