❌

Vue normale

Reçu avant avant-hier

Polars 2.0 pre-release comes with a 5x speed boost — but it could change row order

6 septembre 2026 à 15:30

Working with large datasets can lead to slow queries and out-of-memory errors. Polars, an open-source library that developers and data analysts use to clean, combine, and analyze tables of data, promises to ease both problems in its upcoming 2.0 release. But the first release candidate, out last week, comes with a catch: The new default can change the order of returned rows, potentially affecting code that depends on that order.

In announcing the first release candidate for Polars 2.0, the company says that calling collect on any LazyFrame query will now default to the streaming engine. Per Polars, users can expect “massive memory and performance improvements on most queries,” with the streaming engine expected to be “easily 5x faster” in aggregate.

But they need to keep an eye out for changes in row order. 

Move fast — and maybe re-order things? 

Improved memory usage and performance are obvious upgrades for Polars users who rely on the library for data processing and analysis, and it’s the streaming engine that’s bringing it. 

“Streaming engine doesn’t guarantee row-order by default for certain operations.”

With streaming, Polars says it can execute lazy queries in batches, rather than processing all data at once. This way, users can process datasets that don’t fit into available memory. 

But changing how those queries execute could potentially lead to trouble down the line, as the streaming engine can also change the order in which rows are returned. 

As Polars explains, the “streaming engine doesn’t guarantee row-order by default for certain operations.” That includes operations such as join, group_by, and unpivot. 

In its Version 2.0-rc user guide, the company explicitly calls out the migration hazard and underscores its risk in a red “danger” box, acknowledging that the change “may silently impact the results of your pipelines.”

For users whose code expects rows to appear in a certain order, that could create more problems for downstream processes. 

You can enforce row order, but there’s a chance it may cost you some speed

All is not lost, though. If users are working with code that depends on incidental ordering or observable row order, Polars offers guidance on mitigating the migration risk that comes with the new default. 

The change “may silently impact the results of your pipelines.”

There are two main options: Sort explicitly or set maintain_order=True where applicable.

Alternatively, users can keep the in-memory engine as default by setting the engine affinity. 

What else is coming in Polars 2.0 

Making all LazyFrame queries default to the streaming engine isn’t the only change users can expect from Polars 2.0. Per the announcement, the biggest changes in the upcoming release are improved defaults (the streaming engine being the most significant) and a better API. 

In the pre-release post, Polars explains that 2.0 also removes many ambiguous casts.

For example, it directs users to use .str.to_date()/.str.to_datetime() to parse strings to temporal data types. This way, Polars says users get “one obvious way to parse data.” More examples of improvements to strictness are in the migration guide. 

Why the pre-release before the upcoming Polars 2.0? Because Polars says it “[doesn’t] gate new features” and prefers to ship them as soon as they’re ready.

That said, the company assured users there’s more to look forward to for 2.x, hinting at a new IO-plugin design, a faster S3 reader, a cost-based planner, join reordering, and big SQL coverage improvements, among others.

For developers exploring the release candidate now, the takeaway is clear: Better memory and performance are worth getting excited about, but don’t forget to watch that row order.

The post Polars 2.0 pre-release comes with a 5x speed boost — but it could change row order appeared first on The New Stack.

How Yahoo optimizes resources with flexible VMs in Managed Service for Apache Spark

4 septembre 2026 à 18:00

As a global media and technology company connecting hundreds of millions of users to finance, sports, and entertainment platforms, Yahoo operates a massive data infrastructure where analytics workloads must run continuously at high speed. In deadline-driven data environments, relying on fixed virtual machine (VM) configurations creates a brittle system; if a specific machine shape faces a regional capacity constraint, cluster provisioning in Managed Service for Apache Spark (formerly Dataproc) can experience delays and stall critical data pipelines.

Yahoo utilizes flexible VMs in Managed Service for Apache Spark clusters to automatically absorb these resource fluctuations by defining a ranked list of acceptable VM shapes. This allows the system to dynamically search regional zones and maintain pipeline execution without manual intervention. To search for capacity across a region, teams must also enable Auto-Zone placement.

This optimization builds on Yahoo's broader data modernization journey, which involved migrating on-premises Hadoop and big data estates directly to Google Cloud. By transitioning those legacy workloads, the team established a cloud foundation capable of running high-scale batch and streaming analytics with dynamic resource flexibility.

This post provides a technical blueprint for configuring flexible VM instance rankings in Managed Service for Apache Spark to automatically manage capacity constraints and maintain pipeline execution.

Operational trade-offs of static configurations

Configuring clusters with a single, fixed machine type in a specific zone introduces constraints when regional zonal capacity fluctuations occur, potentially impacting cluster provisioning. Rather than manage these capacity variations through custom retry logic or manual intervention, using flexible configurations allows your infrastructure to automatically adapt. By accepting multiple VM shapes and searching across zones in the selected region, flexible configurations help streamline provisioning to better support high-scale analytics workloads.

Rules for configuring flexible clusters

Deploying flexible configurations requires aligning several connected design choices:

  • Enable auto-zone placement: You must pass a region(--region=${REGION}) or an empty zone string (--zone="") so Managed Spark can search for available capacity across the entire region.

  • Maintain core and memory symmetry: If your Managed Spark cluster uses autoscaling, all machine types in your flexible list must share a similar core count and memory size, even if they come from different VM families. A uniform CPU-to-memory ratio across primary and secondary workers prevents performance degradation, as the smallest ratio determines your effective container sizing.

  • Align component properties: Managed Spark calculates system properties based on VM cores and memory. When mixing machine shapes, you may need explicit property overrides to keep YARN and Spark resource allocations aligned with your expected worker behavior.

Two ways flexible VMs support massive workloads

For large-scale data environments, flexible configurations support operations in two ways:

  1. Higher cluster creation success: Instead of failing when a preferred VM type is out of stock, Managed Spark selects from a ranked list to keep provisioning moving.

  2. Better regional resource use: Auto-zone placement searches the entire region to find capacity, which reduces provisioning friction during high-demand periods.

gcloud example

code_block
<ListValue: [StructValue([('code', 'gcloud dataproc clusters create analytics-cluster \\\r\n --region=us-central1 \\\r\n --zone="" \\\r\n --num-workers=10 \\\r\n --master-instance-selection=\'{"machineTypes":["e2-standard-8"],"rank":0}\' \\\r\n --master-instance-selection=\'{"machineTypes":["n2-standard-8"],"rank":1}\' \\\r\n --worker-instance-selection=\'{"machineTypes":["e2-standard-8"],"rank":0}\' \\\r\n --worker-instance-selection=\'{"machineTypes":["n2-standard-8"],"rank":1}'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7f8e31ac9450>)])]>

API example

You can also build this capacity policy into your automated pipelines or Managed Service for Apache Airflow DAGS using the instanceFlexibilityPolicy field in the ‘Dataproc’ API:

code_block
<ListValue: [StructValue([('code', '{\r\n "projectId": "PROJECT_ID",\r\n "clusterName": "analytics-cluster",\r\n "config": {\r\n "gceClusterConfig": {\r\n "zoneUri": ""\r\n },\r\n "secondaryWorkerConfig": {\r\n "numInstances": 8,\r\n "instanceFlexibilityPolicy": {\r\n "instanceSelectionList": [\r\n {\r\n "machineTypes": ["n2-standard-8"],\r\n "rank": 0\r\n },\r\n {\r\n "machineTypes": ["e2-standard-8", "t2d-standard-8"],\r\n "rank": 1\r\n }\r\n ]\r\n }\r\n }\r\n }\r\n}'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7f8e32253250>)])]>

This API policy achieves the same goal: it establishes your preferred shape, documents valid fallbacks, and lets Managed Spark resolve resource constraints without breaking your automation scripts.

Establishing an infrastructure policy

Managing data at this scale requires standardizing a clear resource policy rather than relying on a single rigid machine type. Your configuration standards should outline:

  • Preferred and fallback VM families for secondary workers.

  • Default auto-zone placement to enable flexible provisioning.

  • Identical core and memory configurations when using autoscaling.

  • Uniform CPU-to-memory ratios across all worker groups to maintain predictable container sizing.

  • Explicit YARN or Spark property overrides to guarantee consistent runtime behavior across different machine lines.

  • Shuffle-safe patterns for Spark workloads running on Spot or highly elastic capacity.

By adopting flexible configurations, you turn infrastructure scarcity into a predictable fallback plan, keeping your critical data pipelines up and running.

Yahoo impact and results

By implementing flexible VMs in Managed Service for Apache Spark, Yahoo successfully reduced cluster provisioning failures by 85% which were caused by regional capacity stockouts. This flexible configuration allows their data infrastructure to automatically handle capacity constraints and successfully provision resources without requiring manual intervention. As a result, Yahoo ensures continuous workload execution and prevents downstream processing delays across their massive data pipelines.

"Managing high-scale data analytics at Yahoo requires resilient, automated infrastructure. Moving to flexible VMs in Managed Service for Apache Spark has transformed our approach; instead of stalling when a specific machine shape faces capacity constraints, our clusters now automatically pivot to our ranked fallback options. This has helped us reduce provisioning failures by 85%, providing the reliability we need to keep our global media platforms running smoothly." - Akshay Jain, Senior Software Developer Engineer, Yahoo!

Strategic benefits of flexible infrastructure

Adopting a flexible compute stack transforms your environment into a dynamic pool of resources that adapts to your operational needs. By moving away from rigid, single-machine type configurations, you ensure that your workloads reliably access the compute they need, regardless of supply fluctuations. This shift not only maximizes workload obtainability and reliability but also facilitates seamless hardware modernization by allowing you to prioritize newer VM generations while maintaining older types as reliable fallback options.

Build your resilient data pipeline

Transitioning to a fluid compute strategy ensures your critical analytics remain operational despite regional resource shifts. Here is how you can begin optimizing your infrastructure today:

  1. Audit your workloads: Identify applications tightly coupled to specific VM families or zones and map out viable alternative hardware shapes.

  2. Standardize resource policies: Explore the documentation for Managed Spark flexible VMs to establish your preferred and fallback VM families.

  3. Align financial strategy: Utilize Flexible Committed Use Discounts (Flex CUDs) to maintain cost predictability when workloads dynamically pivot to alternative machine types.

  4. Claim your credits: New customers may be eligible for $300 in credits to try Managed Service for Apache Spark and other Google Cloud products at no cost.

Serverless Apache Spark on Google Cloud: Architecture Choices & AI Troubleshooting

19 août 2026 à 18:00

In modern enterprise data engineering, Apache Spark remains a cornerstone framework for processing massive datasets at scale. However, managing infrastructure such as provisioning clusters, tuning YARN configurations, and avoiding costs for idle hardware often detracts from what matters most: building resilient data pipelines. Google Cloud addresses this operational overhead via its Managed Service for Apache Spark, offering flexible deployment modes of serverless and managed clusters tailored to specific operational needs.

This technical guide walks through the architectural decision matrix for deploying Spark on Google Cloud, details resource and cost optimization techniques, and demonstrates how to apply built-in Gemini Cloud Assist to rapidly troubleshoot and resolve serverless batch pipeline failures. While there is benefit to reading these three parts in a sequence, each one can be read independently and add value to how you approach Spark development on Google Cloud. 

Part 1: Choosing your Apache Spark deployment model

When launching Spark workloads on Managed Service for Apache Spark, the first major decision point is evaluating whether to construct traditional managed clusters or transition to a zero-management, serverless infrastructure footprint.

Decision #1: Managed clusters vs. serverless

1

*Created using Nano Banana 2 in Gemini Enterprise Agent Platform

Choosing between traditional Managed Spark clusters and serverless depends on ecosystem requirements, infrastructure control needs, and financial utilization patterns:

  • Workload frequency, latency sensitive workloads & financial fit: For continuous, highly predictable, 24/7 streaming or batch processing pipelines where cluster nodes maintain constant high utilization baselines (80%+) or when the workflow’s accumulated startup time risk meeting SLA target, a permanently running, finely tuned traditional cluster, with custom YARN autoscaling rules, can sometimes be more cost-predictable. Conversely, for intermittent, bursty, ad-hoc, or orchestrator-triggered pipelines, Managed Spark serverless is highly optimal, eliminating operational management, requiring less planning time and ensuring you don’t pay for idle compute time.

  • Ecosystem & component requirements: Managed Spark serverless is strictly optimized for Apache Spark 3.x+ codebases. If your processing pipeline relies on other ecosystem components such as Apache Flink, Presto/Trino, Hive LLAP, or Apache HBase, or if you are locked into a legacy Spark 2.x codebase, you must use Managed Spark clusters.

  • Infrastructure customization needs: Managed Spark serverless abstracts away the underlying virtual machine (VM) layer. If your workload mandates deep OS-level hardware tuning, custom OS initialization actions, root SSH access to instances, specific local SSD configurations, or custom machine shapes, a traditional cluster is required. Note that serverless does support custom Docker container images for bundling specific application-level libraries.

Decision #2: Serverless interactive sessions vs. serverless batches

2

*Created using Nano Banana 2 in Gemini Enterprise Agent Platform

Once you select the serverless deployment mode, you must choose the appropriate execution model based on your development stage and operational requirements. Managed Service for Apache Spark provides two options for running serverless workloads:

Serverless interactive sessions

Interactive sessions are great for iterative and exploratory use cases. You write blocks of code, inspect intermediate DataFrames, modify variables, and generate visualizations with your dataset held warm in-memory.

  • Primary interface: Designed for human-in-the-loop interaction. Developers execute code cell-by-cell using their IDE of choice, such as Colab, Gemini Enterprise Agent Platform Workbench, Antigravity, Jupyter notebooks, etc.

  • Idle cost profile: Compute resources remain active to support immediate execution during developer thinking time, which can incur some idle compute charges if sessions are left inactive.

Serverless batches

Batches are useful when you know what you want to run, and need automated, non-interactive execution. The engine runs fully completed, packaged PySpark scripts (.py) or Java/Scala application files (.jar) from start to finish without manual human intervention.

  • Primary interface: Managed by automated orchestrators, such as Managed Service for Apache Airflow, Cloud Scheduler, or CI/CD pipelines.

  • Idle cost profile: Billed strictly for the duration of the run. Compute resources are provisioned on-demand, run the script, and immediately shut down upon completion to prevent idle costs.

The development-to-production lifecycle

These execution options are designed to work together as a natural pipeline lifecycle. During the initial development phase, you open a serverless interactive session within your notebook interface to explore datasets, clean schemas, and prototype transformations. Once your logic is validated and the transformations are finalized, you package the code into a Python script and schedule it as a serverless batch job orchestrated by Managed Service for Apache Airflow for production execution. This transition minimizes ongoing development costs while maintaining operational reliability.

Part 2: Advanced performance tuning and DCU cost optimization

While serverless Managed Spark eliminates the operational overhead of cluster maintenance, running production enterprise-grade pipelines on default settings can result in performance bottlenecks or budget waste. Resource allocation must be explicitly declared during submission using runtime configuration properties to maintain an efficient Data Compute Unit (DCU) burn rate.

Google recently introduced history-based autotuning. In the context of serverless, this capability automatically applies optimizations based on best practices and historical execution. It does this by grouping recurring batch workloads into what Google calls cohorts. The autotuner analyzes the telemetry and statistics from previous runs under that same cohort name to figure out where the bottlenecks are.

Customizing driver and executor shapes

By default, serverless batches allocate generic specifications (4 cores and 16,000MB RAM). This can cause critical efficiency issues depending on the nature of the application:

  • The Memory-Bound job: Pipelines processing highly uncompressed data volumes may hit Out-Of-Memory (OOM) errors and crash. To counter this, increase heap sizing independently using spark.driver.memory and spark.executor.memory.

  • The Compute-Bound Job: Processing-intensive jobs running mathematical modeling or heavy tokenization might saturate CPUs while leaving expensive RAM sitting idle. Fine-tune processing concurrency per instance by explicitly adjusting spark.driver.cores and spark.executor.cores.

Remember that by default increasing cores, automatically provisions a proportionate baseline   of memory to match the vCPU-to-RAM ratio. This is why overriding the values for both cores and memory is critical  

Controlling autoscaling boundaries

Managed Spark serverless dynamically scales up and down the number of active executors based on backlogged tasks. However, unconstrained scaling can lead to budget overruns if a rogue code loop or unoptimized cartesian join is introduced.

As a defensive guardrail, always declare an explicit upper limit using spark.dynamicAllocation.maxExecutors. This acts as your budget deadman-switch. By capping this at a reasonable ceiling, you guarantee that even if the code behaves sub-optimally, the job will never scale past a fixed infrastructure footprint.

  • High priority (SLA-driven): Set maxExecutors to a higher ceiling to allow resource bursting and minimize overall runtime duration.

  • Low priority (nightly batch): Set maxExecutors to a low, tight ceiling. The workload will run longer but will consume a predictable, flat, cost-efficient stream of DCUs.

Managing shuffle storage efficiency

When execution involves wide transformations like groupBy(), join(), or distinct(), data must be redistributed across the network, generating intermediate disk writes known as shuffle storage.

Spark defaults to a static setting of 200 partitions (spark.sql.shuffle.partitions). If you are processing a massive, multi-gigabyte dataset, 200 partitions means each individual chunk will be too large. When a partition's size exceeds available executor RAM (e.g., a 1GB partition trying to process inside 0.5GB of assigned heap space), data spills onto disk. This slows execution and incurs additional billing fees for premium or standard shuffle storage blocks. A helpful rule of thumb: Dynamically scale your partition parameters based on total data size so that each partition handles roughly 100MB to 200MB of data in memory. This may require a few iterations before the optimal results are achieved.

The above properties are the main tunable properties. Additional Serverless runtime configuration properties can be found in this link

Part 3: Operational diagnosis with Gemini Cloud Assist

When automated data pipelines fail in production, data engineers are traditionally forced to spend hours sifting through verbose, disjointed log files across drivers and executors. Managed Service for Apache Spark addresses this friction by natively integrating Gemini Cloud Assist into the Google Cloud console, allowing engineers to diagnose and resolve failures using natural language.

To illustrate this operational shift, we examine the typical troubleshooting lifecycle for a failed PySpark ETL pipeline that reads customer transaction data from a Google Cloud Storage (GCS) bucket, applies transformations, and encounters unexpected runtime errors.

Stage 1: Diagnosing missing execution parameters

During the initial execution attempt of a new pipeline, the batch job status switches from pending to running, and ultimately ends in a failed state with a generic exit message: Application failed with exit code 1.

Rather than manually querying Cloud Logging or navigating through multiple sections of the console, the engineer can locate the error log and select the ‘Investigate log’ option. This action opens a native conversation pane where Gemini Cloud Assist automatically analyzes the driver telemetry and system logs.

3

In this scenario, the assistant explains in plain English that the PySpark script failed because required runtime arguments (such as the source GCS bucket path) were omitted during submission. It instantly identifies the exact lines in the script expecting these arguments, eliminating the need to read through the stack trace.

4

Stage 2: Resolving schema and data type anomalies

Once the missing arguments are resolved and the job is re-submitted, the pipeline runs but encounters a secondary data anomaly. In high-volume ingest pipelines, upstream source files frequently contain corrupted records or formatting inconsistencies.

Upon the second failure, the engineer again prompts Gemini Cloud Assist to investigate the logs. The assistant identifies a TypeError and pinpoints the exact DataFrame transformation causing the crash: a division operation (df['amount'] / df['transaction_id']) that failed because the schema auto-inferred the columns as strings.

5

Additionally, the assistant scans the underlying GCS file data to identify the root cause: non-numeric anomalies (such as text strings within numerical cells) in the source dataset.

6

Stage 3: Generating and deploying verified code fixes

Rather than manually rewriting the PySpark logic to cast schema types and catch null values, the engineer can prompt Gemini Cloud Assist directly to generate a resilient solution:

User Prompt: "Suggest how to rewrite the code to divide the amount by quantity instead of transaction_id. In addition, add logic to skip invalid records without failing the process."

The assistant generates the corrected PySpark code block, using resilient casting and null-handling functions (such as coalesce and try_cast).

By implementing this corrected script, the orchestration pipeline can filter out bad source records smoothly without crashing the entire batch run. The subsequent execution completes successfully, preserving data freshness SLAs.

Unlock serverless Apache Spark: Benefits and next steps

Managing data processing pipelines should not require a deep specialization in infrastructure configuration. By pairing the hands-off scale of serverless batches with explicit resource tuning — such as dynamic allocation caps and calculated shuffle sizing — data teams can maintain strict control over performance and cost profiles. When failures do occur, integrating Gemini Cloud Assist directly into your logging workflows transforms complex troubleshooting from a manual log-sifting exercise into a rapid, automated cycle.

To start putting these architectures into practice, you can explore the Managed Service for Apache Spark documentation and execute a serverless batch directly in the Google Cloud console.

For a deep architectural analysis of these concepts, get instant access to A practitioner’s guide to Apache Spark® in the agentic era. This guide includes step-by-step workflows, Codelabs, and runnable PySpark and Terraform templates directly from our GitHub repository. If you are new to Google Cloud, you can test these blueprints on serverless and managed clusters at zero cost by signing up for a free trial with $300 in credits.

Building cost-effective, high-throughput gen AI workflows in Google Dataflow

18 août 2026 à 18:00

Real-time streaming pipelines are the operational backbone of modern enterprises, continuously processing everything from customer support interactions to transaction logs. Traditionally, streaming DAGs are static; once deployed, their processing logic and execution paths are fixed. However, by integrating generative AI agents, we can move beyond static logic to adaptive execution. This allows streaming workflows to dynamically construct plans, query databases, and trigger custom remediation paths at runtime depending on the content of the data.

For example, when a customer sends an angry message about a damaged order, a pipeline shouldn't just log the error or flag a dashboard. It should look up the order in the database that holds customer order and inventory records, decide on a remediation action (like shipping a replacement or issuing a refund), email the customer, and log the final resolution.

However, streaming systems face a fundamental engineering hurdle when executing gen AI workflows: scale, latency, and cost. Sending every raw event directly to a heavyweight model or multi-step agent equipped with external database and email tools is prohibitively expensive, introduces high latency, and quickly exhausts API rate limits.

This pattern addresses the scale and complexity challenge by combining Google Dataflow, Google Cloud's fully managed, serverless execution service for Apache Beam, and the Agent Development Kit (ADK) to build a hybrid streaming pipeline. By using a lightweight, CPU-bound machine learning model upstream to filter and qualify events, we keep the pipeline highly cost-effective, routing only the complex cases to the downstream agent. There, the agent dynamically decides what actions to take, introducing dynamic branching to the stream without hardcoding thousands of conditional steps into the pipeline's static DAG.

A universal blueprint for high-volume streams

While we use a customer support triage scenario below, this pre-filter + agentic action pattern is a universal paradigm. It applies to any stream where a high volume (>9X%) of events are routine, and only a small number require complex, contextual reasoning.

  • IT Operations & DevOps: Filtering millions of routine system logs on CPU, and triggering an agent to run diagnostics and open bug tickets only when a critical anomaly is flagged.

  • Financial Fraud Triaging: Passing millions of transactions through lightweight, local rules, and calling an agent to execute multi-database lookup tools only for highly suspicious patterns.

  • Industrial IoT: Monitoring normal telemetry on the edge, and routing erratic spikes to an agent to coordinate equipment shutdowns and email field engineers.

The architecture: Why pre-filter streaming events?

In a high-throughput stream, the vast majority of messages do not require complex reasoning or remediation. They might be positive feedback, neutral inquiries, or simple queries.

Routing every single event to a heavyweight LLM workflow creates three primary bottlenecks:

  1. API cost: Frontier models charge per token. Under high throughput, cost scales linearly with stream volume.

  2. Latency: Multi-step workflows (which involve database lookups and external API calls) take seconds, creating a bottleneck in streaming DAGs.

  3. Quotas: External APIs have strict rate limits that streaming workers can easily exhaust.

To prevent this, we build a pre-filtered pipeline in Apache Beam/Dataflow:

image1

Pipeline flow

  1. Ingestion: Read raw customer messages from Google Pub/Sub.

  2. Lightweight sentiment classifier (CPU): Run all messages through a lightweight, CPU-based Hugging Face model (distilbert-base-uncased-finetuned-sst-2-english) using Apache Beam’s RunInference transform. This executes locally on the Dataflow worker CPUs, avoiding external API costs.

  3. Pre-qualification Gate: A simple DoFn filters the stream. Messages with POSITIVE or NEUTRAL sentiment are acknowledged and dropped.

  4. Automated Remediation (ADK): If and only if a message is classified as NEGATIVE, we trigger the gen AI agent backed by gemini-3.5-flash using the ADKAgentModelHandler. The agent uses tools to look up the user in BigQuery, fetch orders, choose a remediation plan, and send a notification email via the Gmail API.

Adaptive execution: Making the Beam DAG dynamic

In traditional streaming architectures, the pipeline's Directed Acyclic Graph (DAG) is rigid. Once deployed to Dataflow, the sequence of transforms is set. If you need to handle new types of alerts or change how specific events are routed, you have to modify, test, and redeploy the entire pipeline.

By placing a gen AI agent downstream of our sentiment pre-filter, we introduce a dynamic, adaptive node inside the static DAG.

For the 95% of records that are positive or neutral, the pipeline runs along a fast, static path. But when the filter gates a negative record, the agent evaluates the payload and dynamically selects the correct sequence of API tools (e.g., database query, inventory check, or email notification) at runtime. This allows the pipeline to execute complex decision trees dynamically, eliminating the need to build and maintain thousands of hardcoded conditional branches in the static Apache Beam code.

Implementing the pipeline

Here is an example implementation in Apache Beam using the Google Agent Development Kit (ADK) and the RunInference framework.

1. Defining the lightweight sentiment model

We define the upstream CPU model using HuggingFacePipelineModelHandler. This model classifies sentiment into POSITIVE, NEUTRAL, or NEGATIVE on the worker instance.

code_block
<ListValue: [StructValue([('code', 'model_handler = HuggingFacePipelineModelHandler(\r\n task="sentiment-analysis",\r\n model="distilbert-base-uncased-finetuned-sst-2-english"\r\n)'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7ffb900e9c10>)])]>

2. Building the heavyweight ADK agent

The ADK agent acts as our remediation assistant. We equip it with three tools:

  • lookup_user: Queries BigQuery for the customer's email.

  • lookup_orders: Queries BigQuery for the customer's orders and current product inventory.

  • send_email: Sends a remediation email to the customer using the Gmail API.

code_block
<ListValue: [StructValue([('code', 'def make_adk_tools(project: str, dataset: str = "sentiment_demo"):\r\n def lookup_user(user_id: int) -> dict:\r\n """Look up user information (email address) from BigQuery by user ID."""\r\n from google.cloud import bigquery\r\n\r\n client = bigquery.Client(project=project)\r\n query = (\r\n f"SELECT user_id, user_email "\r\n f"FROM `{project}.{dataset}.users` "\r\n f"WHERE user_id = @user_id"\r\n )\r\n job_config = bigquery.QueryJobConfig(\r\n query_parameters=[bigquery.ScalarQueryParameter("user_id", "INT64", user_id)]\r\n )\r\n try:\r\n results = list(client.query(query, job_config=job_config).result())\r\n if results:\r\n row = results[0]\r\n return {"user_id": row.user_id, "user_email": row.user_email}\r\n return {"error": f"No user found with user_id={user_id}"}\r\n except Exception as exc:\r\n return {"error": str(exc)}\r\n\r\n def lookup_orders(user_id: int) -> dict:\r\n """Look up a user\'s orders and current product inventory from BigQuery."""\r\n from google.cloud import bigquery\r\n\r\n client = bigquery.Client(project=project)\r\n query = (\r\n f"SELECT p.order_id, p.product_id, pr.remaining_inventory, pr.price "\r\n f"FROM `{project}.{dataset}.purchases` p "\r\n f"JOIN `{project}.{dataset}.products` pr ON p.product_id = pr.product_id "\r\n f"WHERE p.user_id = @user_id"\r\n )\r\n job_config = bigquery.QueryJobConfig(\r\n query_parameters=[bigquery.ScalarQueryParameter("user_id", "INT64", user_id)]\r\n )\r\n try:\r\n results = list(client.query(query, job_config=job_config).result())\r\n orders = [\r\n {\r\n "order_id": row.order_id,\r\n "product_id": row.product_id,\r\n "remaining_inventory": row.remaining_inventory,\r\n "price": float(row.price),\r\n }\r\n for row in results\r\n ]\r\n return {"orders": orders}\r\n except Exception as exc:\r\n return {"error": str(exc)}\r\n\r\n def send_email(to_address: str, subject: str, body: str) -> str:\r\n """Send a plain-text email to the customer via the Gmail API."""\r\n import google.auth\r\n import googleapiclient.discovery\r\n import email.mime.text\r\n import base64\r\n\r\n try:\r\n creds, _ = google.auth.default(\r\n scopes=["https://www.googleapis.com/auth/gmail.send"]\r\n )\r\n service = googleapiclient.discovery.build("gmail", "v1", credentials=creds)\r\n\r\n mime_msg = email.mime.text.MIMEText(body)\r\n mime_msg["to"] = to_address\r\n mime_msg["subject"] = subject\r\n raw = base64.urlsafe_b64encode(mime_msg.as_bytes()).decode("utf-8")\r\n service.users().messages().send(userId="me", body={"raw": raw}).execute()\r\n return "Email sent successfully"\r\n except Exception as exc:\r\n return f"Failed to send email: {exc}"\r\n\r\n return [lookup_user, lookup_orders, send_email]'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7ffb8384c190>)])]>

We configure the LlmAgent and package it in the ADKAgentModelHandler:

code_block
<ListValue: [StructValue([('code', 'adk_agent = LlmAgent(\r\n name="remediation_agent",\r\n model="gemini-3.5-flash",\r\n instruction=(\r\n "You are a customer service remediation assistant with access to "\r\n "BigQuery lookup tools and an email sending tool. "\r\n "When given a prompt describing a customer situation, follow the "\r\n "numbered steps exactly and use your tools to complete the task."\r\n ),\r\n tools=adk_tools,\r\n)\r\n\r\n# RunInference handler for the ADK agent\r\nadk_handler = ADKAgentModelHandler(agent=adk_agent)'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7ffb83940c50>)])]>

3. Assembling the Dataflow DAG

The entire pipeline is declared cleanly. The upstream sentiment inference feeds directly into the filtering step (FilterNegativeADK), which then conditionally executes the downstream ADKInference:

code_block
<ListValue: [StructValue([('code', 'with beam.Pipeline(options=pipeline_options) as p:\r\n # 1. Read from Pub/Sub and classify sentiment on CPU\r\n sentiment_results = (\r\n p\r\n | "ReadFromPubSub" >> beam.io.ReadFromPubSub(topic=known_args.input_topic)\r\n | "DecodeMessages" >> beam.Map(lambda x: x.decode(\'utf-8\'))\r\n | "SentimentInference" >> RunInference(model_handler)\r\n )\r\n\r\n # 2. Filter out non-negative sentiment and invoke the ADK Agent\r\n _ = (\r\n sentiment_results\r\n | "FilterNegativeADK" >> beam.ParDo(FilterNegativeAndPromptADK())\r\n | "ADKInference" >> RunInference(adk_handler)\r\n | "LogADKResults" >> beam.ParDo(LogADKResponse())\r\n )'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7ffb835e8f50>)])]>

Cost and performance advantages

By introducing this filtering step, we gain major engineering and operational advantages:

1. Significant cost reductions

Instead of paying for Gemini input/output tokens on 100% of incoming events, we pay only for the fraction that represent negative customer sentiment (typically < 5% of messages). The other 95% are classified locally on CPU instances at zero incremental API cost.

2. High streaming throughput

Dataflow distributes the CPU classification workload across many instances. Since CPU inference takes milliseconds, the pipeline scales horizontally to handle high-throughput event streams. The heavyweight LLM agent, which can take seconds per request due to tool execution, is called sparingly, preventing backlog.

3. Native Apache Beam integration

Adding the agent into the DAG requires no complex orchestration logic or manual thread pools. Using ADKAgentModelHandler with Beam's native RunInference transform handles parallel worker threads, batching, and integration automatically, keeping the codebase maintainable and clean.

Key takeaways

Streaming data is fast and high-volume, while heavyweight generative AI reasoning is slow and costly.

By building a pre-filtered pipeline with Google Dataflow and the ADK, you get the best of both worlds: the cost and speed of local CPU-based models, and the deep, automated capabilities of Gemini-backed agents.

To see the complete codebase and deploy this yourself, check out the next-2026-demo GitHub repository.


Apache Beam is a trademark of the Apache Software Foundation

❌