Back to Blogs
#ai#googlecloud#jev#dataflow

Event-Driven AI at Cloud Scale: High-Throughput Stream Classification with Pub/Sub, Dataflow, and Jev

How to turn millions of events per hour into typed, confidence-scored decisions, without a...

Updated 8 min read
Cover for Event-Driven AI at Cloud Scale: High-Throughput Stream Classification with Pub/Sub, Dataflow, and Jev

How to turn millions of events per hour into typed, confidence-scored decisions, without a text-generation model in the loop.

In the world of real-time streaming analytics, introducing AI into the critical path is often a recipe for architectural failure.

Processing hundreds of thousands of events per second across clickstreams, transactional audit trails, IoT telemetries, or real-time security logs requires an infrastructure that can absorb ingestion spikes and maintain stable, sub-50ms latencies.

Unfortunately, traditional "System Two" large language models (LLMs) are notorious for token-generation overhead, catastrophic API costs at scale, and non-deterministic JSON outputs that shatter downstream schemas. If you try to call an LLM for every event in a fast pipeline, you will choke your workers and balloon your Google Cloud bill.

TypeSafe AI’s Jev bypasses these bottlenecks. By executing deterministic, microsecond-latency scoring and strict primitive classifications (e.g., boolean noul flags, categorical choice routing, or bounded score), Jev operates like a compiled native function rather than a chat completion endpoint.

Integrating Jev with Google Cloud Pub/Sub and Cloud Dataflow (Apache Beam) creates a stream-native classification architecture capable of enterprise-scale inference without sacrificing throughput.

The problem: classification is a decision, not a conversation
Most event streams don't need a model that writes. They need a model that decides: which team gets this ticket, is this transaction risky, should this log line page someone. Using a chat-style LLM for that means generating text, parsing it, validating it, and hoping it never goes off the rails, at several seconds per call.

This post builds a pipeline around a different kind of model. Pub/Sub buffers events, Dataflow (Apache Beam) scales the processing, and Jev, TypeSafe's first System One model, turns each event into typed answers with calibrated probabilities and confidence scores. Your code, not the model, makes the final call.

What Jev changes

According to TypeSafe's launch announcement, Jev trades string generation for speed and structure. The properties that matter for streaming:
• Typed outputs only. You define the questions in advance; Jev returns a choice, a score, or a yes/no probability. There is no free text to parse, so no parse failures and no off-schema answers.
• Parallel questions in one call. Every question is evaluated independently against the same state in a single request, so adding questions barely changes response time

• Confidence on every answer. Choice and Score answers include probabilities and a confidence value you can gate on.
• Speed and price. TypeSafe reports 70 to 500 ms end-to-end responses and $0.042 per million input tokens, with output tokens free. These are vendor figures from the launch post; measure them on your own data and region before you plan capacity.

• Pub/Sub decouples producers from consumers and absorbs bursts in the subscription backlog.
•** Dataflow** reads the backlog, controls concurrency, and autoscales workers.
• Jev answers the questions you define for each event.
• Your routing code inside the pipeline turns answers and confidence into actions.
• BigQuery and a routing topic store decisions and fan them out. A dead-letter topic catches permanent failures.

Backpressure is built in: if Jev slows or throttles, Dataflow slows down and the backlog grows instead of events being dropped. Pub/Sub delivers at least once, so key your sink on a stable event ID to make replays harmless.

Step 1: Define atomic questions

Jev works best when each question asks one narrow thing, the kind of gut-check a knowledgeable person could answer in seconds. Decompose a fuzzy judgment into several questions and combine them in code .For support tickets:

from typesafe_sdk import Choice, Noul, Score

QUESTIONS = {
    'department': Choice(
        instructions='Which team should handle this ticket?',
        criteria={
            'billing': 'Payment, invoice or subscription issues',
            'technical': 'Bugs or integration problems',
            'account': 'Login, access or account settings',
            'sales': 'Pricing or plan questions',
            'other': 'Anything else',
        },
    ),
    'urgency': Score(
        instructions='How urgent is this ticket?',
        criteria=[
            'Routine, no time pressure',
            'Time-sensitive',
            'Critical: outage or lost revenue',
        ],
    ),
    'is_abuse': Noul(
        instructions='The message contains abuse, threats or harassment',
    ),
}
Enter fullscreen mode Exit fullscreen mode

Three question types cover most classification work: Choice picks one option from a closed set, Score rates against ordered levels, and Noul returns the probability that a statement is true. All three can be mixed in one call, and the label set lives in your code, so a typo in a department name is a code review problem rather than a production incident

Step 2: The streaming pipeline

Jev evaluates one state per request, so the throughput lever is concurrency rather than packing many events into one prompt. The pipeline below groups events into small batches only to fan out concurrent async calls per worker thread.

import asyncio
import json
from dataclasses import asdict, dataclass

import apache_beam as beam
from apache_beam.io.gcp.pubsub import ReadFromPubSub, WriteToPubSub
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.transforms.util import BatchElements
from typesafe_sdk import AsyncTypeSafeClient, RetryPolicy


@dataclass(frozen=True)
class Decision:
    id: str
    department: str
    department_confidence: float
    urgency: float
    is_abuse: float
    route: str


def route(department: str, confidence: float, urgency: float, abuse: float) -> str:
    if abuse >= 0.8:
        return 'trust_and_safety'
    if confidence < 0.6:
        return 'human_triage'
    if urgency >= 2 and confidence >= 0.85:
        return 'page_oncall'
    return department


class ClassifyFn(beam.DoFn):
    def setup(self):
        self.retry = RetryPolicy(max_retries=3, backoff_max=1.0, timeout=5.0)

    async def _run(self, batch):
        async with AsyncTypeSafeClient(retry=self.retry) as client:
            return await asyncio.gather(
                *(client.system_one(ev['text'][:4000], QUESTIONS) for ev in batch),
                return_exceptions=True,
            )

    def process(self, batch):
        results = asyncio.run(self._run(batch))
        for ev, res in zip(batch, results):
            if isinstance(res, Exception):
                yield beam.pvalue.TaggedOutput(
                    'failed', {'id': ev['id'], 'payload': ev, 'error': str(res)[:300]})
                continue
            dept = res.choices['department']
            urgency = res.scores['urgency'].score
            abuse = res.nouls['is_abuse'].noul
            yield Decision(
                id=ev['id'],
                department=dept.choice,
                department_confidence=dept.confidence,
                urgency=urgency,
                is_abuse=abuse,
                route=route(dept.choice, dept.confidence, urgency, abuse),
            )


def run():
    opts = PipelineOptions(streaming=True, save_main_session=True)
    with beam.Pipeline(options=opts) as p:
        out = (
            p
            | 'Read' >> ReadFromPubSub(
                subscription='projects/my-project/subscriptions/tickets-sub')
            | 'Parse' >> beam.Map(lambda b: json.loads(b.decode('utf-8')))
            | 'Window' >> beam.WindowInto(beam.window.FixedWindows(10))
            | 'Batch' >> BatchElements(min_batch_size=16, max_batch_size=64)
            | 'Classify' >> beam.ParDo(ClassifyFn()).with_outputs('failed', main='ok')
        )

        decisions = out.ok | 'ToDict' >> beam.Map(asdict)
        decisions | 'ToBQ' >> beam.io.WriteToBigQuery(
            'my-project:ai.ticket_decisions',
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
        decisions | 'ToJson' >> beam.Map(lambda d: json.dumps(d).encode()) \
                  | 'ToRouting' >> WriteToPubSub('projects/my-project/topics/ticket-routes')

        out.failed | 'DLQJson' >> beam.Map(lambda d: json.dumps(d).encode()) \
                   | 'ToDLQ' >> WriteToPubSub('projects/my-project/topics/tickets-dlq')
Enter fullscreen mode Exit fullscreen mode

The SDK calls are the documented system_one method on AsyncTypeSafeClient, and the client defaults to jev-latest and reads TYPESAFE_API_KEY from the environment. Treat the Beam wiring as a starting point and adjust it to your Beam and SDK versions.
Notes on the design:

• The Decision dataclass is frozen, so every downstream stage sees one immutable, fully typed shape.
• return_exceptions=True means one failed call becomes one dead-letter record instead of a retry storm for the whole batch.
• The client is opened per batch inside asyncio.run, which keeps its connections tied to the event loop that uses them.
• Truncating text bounds token usage and guards against pathological payloads.

Step 3: Gate on confidence, not just on the label

The route function above is the heart of the design. The answer says what; confidence says whether to act (confidence-gated routing pattern). Each action gets a threshold that matches the cost of being wrong:

• Below 0.6 confidence, send the ticket to a human triage queue.
• Abuse at or above 0.8 goes to trust and safety.
• Paging on-call requires both top urgency and 0.85 confidence, because a false page is expensive.
• Everything else routes to its department automatically.
These numbers are starting points. Calibrate them on a labelled sample from your own traffic, then move the thresholds in code without touching a prompt.

Step 4: Size concurrency and respect limits

Throughput follows from concurrency and latency:

sustained events/sec = concurrent calls / average call latency
Enter fullscreen mode Exit fullscreen mode

If you allow 200 concurrent calls and the average latency is 0.3 seconds, you sustain about 670 events per second, or roughly 2.4 million events per hour. Control it with three levers:

  1. Cap max_num_workers and threads per worker so total in-flight calls stay under your account's limits. Check your plan's rate limits with TypeSafe; they are not covered here.
  2. Tune RetryPolicy (retry count, backoff, timeout) so transient errors retry briefly and then dead-letter.
  3. Enable HTTP/2 with the typesafe-sdk[http2] extra, which the docs recommend for many concurrent requests because it multiplexes calls over one connection.

Step 5: Make failures boring

• Permanent failures go to the dead-letter topic with the original payload and error. Build a small replay job to drain it once the cause is fixed.
• Duplicates are handled by keying BigQuery on the event ID: write to a staging table and MERGE, or dedupe at query time.
• Outages or throttling. The backlog simply grows. For critical events, add a rules-based fallback that routes to human triage if Jev is unavailable for longer than your latency budget.

Step 6: Control cost

Because Jev bills input tokens only, cost scales with how much text you send. TypeSafe's quick-start example used 392 input tokens for a short ticket with three questions (docs). At roughly 400 tokens per event and $0.042 per million tokens, one million events is about 400 million tokens, or around $17 in Jev input charges. That is an estimate from published list pricing, so confirm it against your own token counts.
To keep it low:
• Strip signatures, quoted replies and boilerplate before classification.
• Cache results for exact repeats, which are common in logs and templated messages.
• Ask only the questions you actually act on; speculative questions are cheap but not free.

Step 7: Observe everything

Track these signals:

• Pub/Sub: oldest unacknowledged message age (your end-to-end lag) and backlog size.

• Dataflow: system lag, worker count and CPU.

• Jev calls: latency percentiles, error rate by status, tokens per event.

• Decision quality: route mix over time, average confidence, and the share sent to human_triage. A sudden shift usually means upstream data changed.

Sample a small percentage of automatic decisions for human review. If high-confidence answers are wrong more often than your thresholds assume, tighten the gates.

When to use Jev, and when not to

Jev gives up text generation by design, so it classifies, scores, routes and extracts, but it does not write replies or summaries. A common split is to let Jev decide at stream speed and call a generative LLM only for the small slice of events that need written output, such as a drafted response to an urgent ticket. TypeSafe also publishes known rough edges for each model version, so read the Jev 1.13 jaggedness notes before trusting it on a new task.

Wrapping up

Pub/Sub gives you a durable buffer, Dataflow gives you elastic processing, and Jev gives you typed, confidence-scored decisions at stream speed. The key design choice is to keep control in your code: small atomic questions in, typed answers out, and thresholds you can change without retraining or reprompting.
Start small: one topic, one subscription, three questions, one table. Measure lag, cost and the confidence distribution, then tune the gates.

Sources

• Introducing System One Models & Jev, TypeSafe AI, Sep 15, 2026
• Jev introduction and Quick start, TypeSafe docs
• Python SDK usage, TypeSafe docs
• Confidence-gated routing, TypeSafe docs