apache-beam-unified-batch-and-stream

apache-beam-unified-batch-and-stream is a skill for Claude Code, Codex from vaquarkhan/data-engineering-agent-skills. It costs 45 tokens per session (976 once invoked), scanned A, original, MIT.

A guide for Apache Beam, a framework for writing data-processing pipelines that can handle both stored data and continuously arriving data. It covers how event time, late data, windows, and execution backends affect results.

In plain words
What is it for?
Use it to design Beam transforms, shared batch-and-streaming logic, windowing, watermarks, triggers, replay behavior, and pipelines that can run on Dataflow, Flink, Spark, or Beam's local runner.
Why use it?
It helps keep batch and streaming pipelines consistent while avoiding incorrect assumptions about time or the system that runs them.

Skill for Claude CodeCodex

Written for no agent in particular: nothing here depends on one.

Good fit Use it to design Beam transforms, shared batch-and-streaming logic, windowing, watermarks, triggers, replay behavior, and pipelines that can run on Dataflow, Flink, Spark, or Beam's local runner.

Compare 6 skills from other repositories ↓
Install with agentmods
npx agentmods add skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream
Install

Getting it into your agent

One page per mod, every tool's command on it. A separate URL per tool would split the same page into five that compete with each other.

Any agent
npx skills add vaquarkhan/data-engineering-agent-skills --skill apache-beam-unified-batch-and-stream
Clone the repo
git clone --depth 1 https://github.com/vaquarkhan/data-engineering-agent-skills

Made for: Claude Code, Codex.

Wrote this? Show the measurements

A badge with what this costs and how it scanned, read live from this page, so it follows the numbers instead of freezing them. Markdown for a README, HTML for a documentation site or a project page.

agentmods badge for apache-beam-unified-batch-and-stream

README.md
[![agentmods](https://agentmods.dev/badge/skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream/github.svg)](https://agentmods.dev/skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream)
Your own site
<a href="https://agentmods.dev/skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream"><img src="https://agentmods.dev/badge/skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream/github.svg" alt="Measured on agentmods" height="20"></a>

Or the 80×15 button, for a site that already has a row of RSS and ATOM ones. Only the verdict fits; the numbers stay here.

agentmods 80×15 button for apache-beam-unified-batch-and-stream

Your own site · 80×15
<a href="https://agentmods.dev/skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream"><img src="https://agentmods.dev/badge/skills/vaquarkhan/data-engineering-agent-skills/apache-beam-unified-batch-and-stream.svg" alt="Reviewed on agentmods" width="80" height="20"></a>
Per session 45 Skills are progressive disclosure: only the name and description are preloaded; the body loads when the skill is used.
When invoked 976 The whole file, excluding the scripts and references it only reads on demand.
Security scan A 0 findings. A grade says what 26 rules found in the file — not that it is safe.
Origin original No closer match found in the catalogue.
Token cost

What it costs to keep this loaded

Counted locally with the o200k_base tokenizer, which is exact for GPT models; Claude uses its own tokenizer and its counts differ. Treat this as one consistent yardstick across the catalogue rather than a bill. Prices are per million input tokens.

ModelPer sessionOnce invoked
Fable 5.1 $0.00045 $0.00976
Opus 5 $0.00023 $0.00488
Sonnet 5 $0.00009 $0.00195
Haiku 4.5 $0.00005 $0.00098

Measured 10d ago against content hash 39fa8666190b, method: parsed. Prices are Anthropic first-party input rates as of 2026-09-10, from the pricing page.

Security

Grade A, and why

apache-beam-unified-batch-and-stream scanned grade A with 0 findings against 26 rules in 11 categories — prompt injection, anti-refusal, data exfiltration, privilege escalation, supply chain, agent snooping, system-prompt leakage, SSRF and excessive agency — measured 10d ago.

A static scan of the body, not an audit. Every finding is printed with the line that produced it so you can judge whether it matters here. A mod is markdown that instructs an agent; that is exactly why what it instructs is worth reading.

Nothing flagged

None of the 26 patterns this scan looks for appear in this file: no shell pipes, no recursive deletes, no credential paths, no hidden text, no instruction-override or anti-refusal phrasing, no agent-config snooping. That is not a guarantee, it is the absence of the things that are checkable.

skills/apache-beam-unified-batch-and-stream/SKILL.md · 90 lines

How it starts

The opening of the file, as written. The whole thing — 90 lines — stays where its author put it; the contents beside it link to each section on GitHub.

Apache Beam Unified Batch And Stream

Overview

Use this skill when Apache Beam is the abstraction layer for both batch and streaming data processing. It helps agents preserve portability without hiding time semantics or runner-specific constraints.

When to Use

  • building Apache Beam pipelines for batch, streaming, or unified workloads
  • targeting multiple runners such as Dataflow, Flink, Spark, or Direct Runner
  • sharing transform logic across batch and streaming modes
  • managing windowing, watermarks, triggers, and late data handling
  • designing portable pipelines that must run across environments

Do not use this when the workload is locked to a single runner and Beam portability is not a goal.

Workflow

  1. Define the pipeline contract and time semantics. Include:

    • input PCollections and their bounded or unbounded nature
    • event-time versus processing-time expectations
    • output schema, grain, and freshness requirements
    • delivery guarantees expected by downstream consumers
  2. Separate portable pipeline logic from runner-specific deployment.

    • keep transforms, DoFns, and combiners runner-agnostic
    • isolate runner configuration (parallelism, autoscaling, resource hints) into pipeline options
    • document which runners are supported and tested
    • avoid runner-specific APIs unless portability is explicitly sacrificed
  3. Design windowing and trigger strategy explicitly. Account for:

    • fixed, sliding, session, or global windows
    • trigger behavior: when to emit, accumulate, or discard
    • allowed lateness and late data routing
    • watermark advancement assumptions per source
  4. Handle state and side inputs carefully.

    • stateful DoFns bind to a specific key space — document key cardinality
    • side inputs can become bottlenecks at scale — prefer bounded and small
    • timers must account for watermark-driven versus processing-time semantics
    • state cleanup must be explicit for unbounded pipelines
  5. Make testing and local validation part of the workflow.

    • use DirectRunner for correctness tests
    • validate windowing behavior with synthetic watermark progression
    • test exactly-once semantics through pipeline drains and restarts
    • confirm output idempotency for at-least-once runners

Read the full file on GitHub · 90 lines

Changes

What this file has done since we first saw it

Hashed on every crawl. A supply-chain change to an agent config is a question of when, not whether, so the history is kept rather than the latest state alone.

  1. 10d ago First seen · 90 lines · 45 tokens per session scan A 39fa8666190b

Subscribe to this mod's changes

apache-beam-unified-batch-and-stream is a skill published in the GitHub repository vaquarkhan/data-engineering-agent-skills (44 stars, last pushed 2mo ago), licensed MIT. It adds 45 tokens to every session and 976 once invoked, about $0.0002 per session on Opus 5. A static security scan graded it A with 0 findings. No closer match exists in the catalogue, so it is treated as the original; first seen 2026-08-30.

Related

Other skills, from other repositories

kafka-shadowtraffic

Generate a ShadowTraffic configuration to populate a Kafka topic with realistic synthetic data. Discovers the target topic, its key and value schemas, and the correct serializers from the live cluster via any attached Kafka MCP server, then writes a ready-to-run shadowtraffic-config.json and Docker command. Use when…

lensesio/agentic-engineering-for-apache-kafka · 134 tokens

claude-api

Build, debug, and optimize Claude API / Anthropic SDK apps. Apps built with this skill should include prompt caching. Also handles migrating existing Claude API code between Claude model versions (4.5 → 4.6, 4.6 → 4.7, retired-model replacements). TRIGGER when: code imports anthropic/@anthropic-ai/sdk; user asks for…

Prismer-AI/PrismerCloud · 193 tokens

migrating-ai-sdk-to-common-ai

Migrates Airflow projects from airflow-ai-sdk to apache-airflow-providers-common-ai 0.4.0+. Use when replacing airflow-ai-sdk with the official Airflow AI provider - migrating LLM decorators (@task.llm, @task.agent, @task.llmbranch, @task.embed), switching from model strings/objects to connection-based LLM…

astronomer/agents · 151 tokens

creating-openlineage-extractors

Create custom OpenLineage extractors for Airflow operators. Use when the user needs lineage from unsupported or third-party operators, wants column-level lineage, or needs complex extraction logic beyond what inlets/outlets provide.

astronomer/agents · 51 tokens

telnyx-ai-inference-curl

Access Telnyx LLM inference APIs, embeddings, and AI analytics for call insights and summaries. This skill provides REST API (curl) examples.

team-telnyx/ai · 39 tokens

kafka-dlq-review

Review dead letter queue implementations for completeness using the Lenses MCP server. Checks DLQ topic existence, configuration, monitoring, metadata preservation, retry logic, reprocessing paths and connector DLQ alignment. Use when user says "review dead letter queues", "check DLQ setup", "DLQ audit" or asks about…

lensesio/agentic-engineering-for-apache-kafka · 93 tokens