GCP News - 2026-08-27

2026-08-27
最終更新: 2026-08-27 21:31:29 JST

Google Cloud Blog

How Uber improves network reliability while unblocking cloud migration

詳細を表示

Uber has a lot in common with the cities it serves. Both are always changing and growing, both must carefully manage the resulting traffic to prevent congestion and sprawl.

Uber has continuously evolved its technical strategies to manage its expanding network, and this careful planning and constant evolution helps ensure that application traffic across its entire platform runs smoothly. Ultimately, maintaining a reliable, high-scale platform that operates seamlessly at any given time is key to preserving user trust.

One important solution in this effort has been application awareness on Cloud Interconnect. An industry-first tool for application prioritization across hybrid networks, application awareness on Cloud Interconnect has helped Uber prioritize critical traffic to ensure business continuity during potential network congestion events. 

Uber acted as an early design partner for application awareness on Cloud Interconnect, helping ensure that this capability met the demands of Uber’s global-scale operations. It not only improved Uber’s daily operations, it also gave Uber the confidence to move forward with a Google Cloud migration, with confidence that there would be less risk of service interruptions during switchovers. 

In this post, we’ll explain the features Uber most sought and why, the inner workings of application awareness on Cloud Interconnect, and how it can help other organizations as well.

Prioritizing critical traffic

When migrating distributed, hybrid, or multicloud applications at a global scale, network reliability becomes a primary concern. Even the most worthwhile migrations may not seem worth it if such migrations interrupt ongoing service. For organizations like Uber, moving vast amounts of data to support large data analytics workload — including emerging AI use cases — can saturate network links, resulting in increased reliability risk for their critical application traffic. 

With standard cloud interconnect approaches, enterprises typically apply simple bandwidth overprovisioning to meet extreme infrastructure needs. But with today's hybrid cloud demands, and given the size of an organization like Uber, overprovisioning network capacity for peak usage is often too costly and unreliable. 

The shortcomings of overprovisioning only become magnified with the integration of cutting-edge AI innovations. Uber needs systems in place that can take on massive data transfers without congesting its network and protecting the performance of business-critical applications.

With the benefit of application awareness on Cloud Interconnect, including the four major features of application awareness — traffic handling, congestion response, latency management, and cost efficiency — Uber was able to achieve the networking optimization its modern tech stack requires.

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="aai concept value prop with_without picture" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/aai_concept_value_prop_with_without_pictur.max-1000x1000.jpg" />
    
    </a>
  
</figure>


  </div>
</div>

Starting with a private preview, Uber deployed this feature across its infrastructure, beginning with Google Cloud Interconnect deployments in Phoenix, Arizona, and Ashburn, Virginia. Application awareness on Cloud Interconnect allows Uber to classify and prioritize end-user application traffic over less time-sensitive data using DSCP marking and configured queuing profiles.

In the following chart, we look at the four key features of application awareness on Cloud Interconnect, how they differ from legacy approaches, and how they help provide better operational continuity for organizations like Uber. 

Feature

Standard interconnect solutions

Application awareness on Cloud Interconnect

Traffic handling

All traffic treated equally (first-in, first-out)

Traffic classified into six distinct traffic classes

Congestion response

High-priority application traffic may be dropped during bursts

Business-critical traffic is protected via strict priority or bandwidth sharing policies

Latency management

Unpredictable latency for high priority applications

Predictable and consistent low-latency for time-sensitive workloads

Cost efficiency

Requires expensive overprovisioning to absorb peaks

Efficient bandwidth utilization and lower TCO

Uber's key takeaways

For Uber, the business value of being able to prioritize business-critical traffic on its networks by deploying application awareness on Cloud Interconnect was immediate. And in doing so, Uber has also created a blueprint that other enterprises with similar hybrid cloud challenges can replicate. The core elements of that blueprint include:

  • Ensuring business continuity: Uber can decide in real time which application traffic to prioritize during major, high-traffic events. This means that mission critical applications stay up and running during even extreme events (both planned and unplanned). Uber leadership has called application awareness on Cloud Interconnect important for its global operations. 

  • Efficient bandwidth utilization: Instead of blindly overprovisioning bandwidth to prevent congestion, application awareness allows Uber to better utilize their existing Cloud Interconnect capacity aligned with their expected network bandwidth needs. The result is lower total cost of ownership for network infrastructure.

  • Unblocked workload migration: By protecting critical applications from network congestion, Uber was able to migrate significant workloads to Google Cloud and, in the process, dramatically reduce operational overhead.

"Application awareness on Cloud Interconnect was the key that unlocked our ability to migrate more strategic workloads to Google Cloud and is critical for maintaining service reliability during peak global demand. By allowing us to intelligently prioritize traffic, it helps us ensure that we can protect our higher priority services and make our infrastructure more efficient, lowering our total cost of ownership. This wasn't just a feature deployment; it was a deep engineering partnership that delivered a solution critical to our business." Harry Liu, Director of Engineering, Uber

Securing network reliability for AI and beyond

As more enterprises integrate cloud-based AI models, distributed applications, and data analytics, it's becoming a business imperative to be ready to handle the massive data transfers that follow. But in doing so, they also have to ensure they never compromise the reliability of their critical applications. 

With application awareness on Cloud Interconnect, Uber demonstrated that moving beyond simple bandwidth overprovisioning to protect business-critical traffic was an essential step to building the stability required to embrace modern hybrid and multicloud strategies.

You can read our blog about the potential of Cloud Interconnect across industries to learn more about what the service can bring to your organization, and if you’re ready to explore more, our team of networking and industry experts are ready to help.

Related Article

    <div class="uni-related-article-tout__content-wrapper">
      <div class="uni-related-article-tout__image-wrapper">
        <div class="uni-related-article-tout__image"></div>
      </div>
      <div class="uni-related-article-tout__content">
        <h4 class="uni-related-article-tout__header h-has-bottom-margin">How Vodafone is using gen AI to enhance network life cycle</h4>
        <p class="uni-related-article-tout__body">Vodafone and Google Cloud deployed generative AI to unlock new levels of efficiency, creativity, and customer satisfaction through networ...</p>
        <div class="cta module-cta h-c-copy  uni-related-article-tout__cta muted">
          <span class="nowrap">Read Article
            <svg class="icon h-c-icon" xmlns="http://www.w3.org/2000/svg">
              <use xlink:href="#mi-arrow-forward" xmlns:xlink="http://www.w3.org/1999/xlink"></use>
            </svg>
          </span>
        </div>
      </div>
    </div>
  </div>
</a>

Simplify your resilience testing strategy with Fault Injection Testing

詳細を表示

When databases fail and network paths falter, you still need your mission-critical cloud services to stay online. Yet guaranteeing high availability has become increasingly difficult because of the complexity of modern distributed systems. 

To help you maintain availability and reliability during adverse events, we’re announcing Fault Injection Testing in preview. Fault Injection Testing is designed to help developers and architects automate failure testing to ensure predictable behavior during disruptions. 

By deliberately introducing faults into your environment, you can verify your safety mechanisms before an actual outage impacts your customers.

Why native resilience testing matters

Unlike in self-hosted data centers, cloud applications offer less direct access to underlying infrastructure to facilitate failover testing.

Without native tools to prove your application can survive a failure, you risk a critical gap in your reliability strategy that exposes you to several risks:

  • Damaged trust and reputation: Frequent failures or poor performance lead to customer dissatisfaction and long-term damage to your brand's image.

  • Compliance and regulatory penalties: For many industries, particularly financial institutions, failing to prove disaster recovery capabilities can lead to non-compliance, audits, and fines.

  • Migration delays: Large-scale migrations often stop when teams cannot verify that critical applications will remain stable during a zone failure.

How Fault Injection Testing works

Fault Injection Testing allows you to run experiments by creating experiment templates. These templates act as blueprints, defining the specific fault to be injected and the resources that will be targeted for the experiment.

In this public preview, you can test two primary failure scenarios:

  • Failover Cloud SQL: This fault triggers a failover of a high availability Cloud SQL instance from the primary zone to a standby zone.

  • Degrade application traffic: This allows you to selectively add latency and HTTP error codes through a Layer 7 load balancer.

Before any fault is injected, Fault Injection Testing performs an automated dry run. This read-only simulation checks your permissions and provides an up-to-date list of every resource that will be affected. 

Once you verify the scope, you can manually start the injection. The duration you defined in the template will run its course, and the faults will be reverted at the expiration of the timer.  

During the experiment, you can verify that your application is behaving as you planned.  If things do not go as planned, you can use the stop and revert capability to immediately halt the experiment and begin restoring resources to their normal state.

During preview, we recommend as a best practice to use FIT in a non-production environment. Preview is an opportunity to get early access to learn how the service fits and complements your existing testing practices, and to provide us with your feedback to improve the product as well!

Built for the enterprise

Partners like KeyBank and Servier are already using Fault Injection Testing to validate their deployments. By using native fault injection, these organizations can approximate demanding failure scenarios — such as zonal outages — to help ensure their services remain stable.

Get started with Fault Injection Testing

Fault Injection Testing is available through the Google Cloud console, the gcloud CLI, and REST APIs.

  1. Request preview access: Talk to your Google Cloud Account Team to add your project to the preview.

  2. Enable the API: Search for "Fault Testing API" in your Google Cloud console and select enable.

  3. Assign roles: Ensure your team has the roles/faulttesting.operator role to configure and run experiments.

  4. Run your first dry run: Create a template for a Cloud SQL or load balancer resource in a non-production environment and execute a dry run to see the potential impact.

For more details on implementation, talk to your account team, or view the User Guide for Fault Injection Testing.

Using OKF with Knowledge Catalog to serve context for agents

詳細を表示

We continue to iterate on the Open Knowledge Format (OKF), an open specification that formalizes the LLM-wiki pattern into a portable, interoperable format. But a big question remains: How can you share and govern access to an OKF bundle across an organization?

OKF v0.1 established a portable format for the context agents need: markdown files with YAML frontmatter, one required field, and five conventions. Then, OKF v0.2 added the trust signals (provenance, verification, freshness, attestation) that a machine-authored bundle requires to be relied on, allowing a team to publish a trustworthy bundle for its own agents. 

However, what OKF does not answer is how teams share their bundles across an organization. A git repo per bundle is portable, but it is not searchable alongside the data it describes, it cannot be secured and governed using the same organizational identity and compliance policies, and it does not sit next to the technical metadata (schemas, lineage, ownership) that data teams already work in. Every downstream agent must know where each bundle resides, and that does not scale beyond a small number of bundles.

To scale an OKF bundle across an organization, you can use Knowledge Catalog, Google Cloud's context engine for agents. By mapping the bundle onto Knowledge Catalog's existing types, every concept becomes discoverable, governed, and reachable by any agent already reading from the catalog.

Knowledge Catalog is the context engine for agents

Every agent that queries Knowledge Catalog reads from one governed index over what the organization already has in BigQuery, Cloud Storage, operational databases, and applications. Each entry carries schema, lineage, ownership, and tags, and can be extended with typed aspects that add domain-specific fields. The same catalog exposes search and cross-project lookup to retrieve optimized context for each agentic query. The context retrieval is secure and governed by IAM controls, so agents can only see the entries they have access to based on IAM identity. 

Publishing an OKF bundle into Knowledge Catalog takes a one-time setup and a single push. Both use the OKF sample code in the Knowledge Catalog repository, whose wrappers call gcloud dataplex for setup and delegate push to kcmd (the Metadata-as-Code CLI in the same repository).

The setup registers three Knowledge Catalog resources: an EntryGroup to hold the bundle, an EntryType named okf-bundle for its concepts, and an AspectType named okf that carries the OKF signal fields (from the okf-aspect.json schema in the sample code). The push then creates one okf-bundle Entry per concept, each with two Aspects: an overview Aspect for the markdown body, and an okf Aspect for the structured signals. Display name, description, and tags live on the Entry itself. The bundle's index.md navigation files and its root log.md are also published as Entries: index files carry only the overview Aspect (no OKF frontmatter), and log.md carries both Aspects with okf_type: Log.

Everything Knowledge Catalog already does for technical metadata (search, IAM, lineage, cross-project discovery) applies equally to OKF bundles, alongside the data they describe.

The okf AspectType

The okf-aspect.json schema in the sample code defines the AspectType. It carries 13 fields covering the full OKF v0.2 spec:

#

Field

Type

Purpose

1

okf_type

string

The OKF document type (freeform, e.g. BigQuery Table, Metric, Attested Computation).

2

generated

record {by, at}

Actor and timestamp for the last meaningful change.

3

sources

array of {id, resource, title, author, usage_count, last_modified}

Materials the concept derives from, with credibility signals.

4

verified

array of {by, at}

Verification events. A human: actor marks the highest trust tier.

5

status

string

Lifecycle state: draft, stable, or deprecated.

6

stale_after

datetime

Absolute point in time (RFC3339 with an explicit offset) on or after which the content is stale.

7

usage_window

record {from, to}

Period the source usage counts were measured over.

8

runtime

string

How an Attested Computation runs (e.g., bigquery).

9

parameters

array of {name, type, required}

Typed named holes a caller may fill. The only surface a caller may vary.

10

computation

string

Path to a file holding the computation body.

11

executor

record {resource, receipt[]}

How the computation runs and what evidence it must return.

12

attester

record {resource}

Deterministic code that takes a receipt and returns a verdict.

13

extra

string

Producer-defined frontmatter the template does not model, as JSON [path, value] pairs. Keeps the round-trip lossless.

Every field is annotated with a display name, a description, and a mandatory index. Any top-level scalar field in the okf Aspect (okf_type, status, stale_after, runtime, computation, extra) can drive Knowledge Catalog search predicates directly, so aspect:acme-analytics.us-central1.okf.okf_type=Metric returns every OKF Metric in scope. Scalar subfields of record fields (generated.by, usage_window.from, executor.resource, attester.resource) also drive predicates. The array fields (sources, verified, parameters) are not server-side searchable on their subfields; agents narrow on them client-side after entries.get with view=ALL. One caveat for search predicates on datetime-typed fields (stale_after, generated.at, usage_window.from/.to), use a bare date (stale_after=2026-12-31) or a range comparison (stale_after>2026-01-01), not the full RFC3339 timestamp.

Pushing a bundle

kcmd push reads an OKF bundle from git and writes each concept as an Entry in the target Knowledge Catalog EntryGroup. index.md files become Entries too, and each concept is parented to the index above it, so the bundle's directory structure survives as a browsable hierarchy.

kcmd expects a bundle in the Documents Layout: markdown files under a catalog/ subdirectory, and a catalog.yaml at the bundle root that lists the snapshot's entry and aspect types. The sample code's setup.ts generates catalog.yaml from its --entry-group flag (default okf_demo), so a reader wiring the sample to a new bundle passes the flag rather than editing catalog.yaml by hand.

Here is an end-to-end workflow for the Acme Retail bundle that we introduced in the OKF v0.2 blog:

code_block
<ListValue: [StructValue([('code', '# One-time setup (if required): install bun, clone the repo, build kcmd, configure gcloud\r\ncurl -fsSL https://bun.sh/install | bash\r\nexport BUN_INSTALL="$HOME/.bun" && export PATH="$BUN_INSTALL/bin:$PATH"\r\ngit clone https://github.com/GoogleCloudPlatform/knowledge-catalog\r\ncd knowledge-catalog/toolbox/mdcode && npm install && npm run build\r\n\r\n# Authenticate, set project and enable dataplex apis\r\ngcloud auth login\r\ngcloud config set project <your-project>\r\ngcloud config set compute/region <your-location>\r\ngcloud services enable dataplex.googleapis.com\r\ngcloud auth application-default login\r\n\r\n# Push the Acme Retail bundle\r\ncd demo/okf\r\nbun run setup.ts # creates the EG (default \'okf_demo\')\r\nbun run push.ts # pushes okf/bundles/acme_retail into the EG setup created'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7f590bb508d0>)])]>

To pick a different EntryGroup name or push a different bundle, pass --entry-group your-name to setup.ts and --bundle path/to/your/bundle to push.ts. For example: bun run setup.ts --entry-group acme-bundle followed by bun run push.ts. This regenerates the manifest, so subsequent push, pull, and cleanup all target the new EG; delete earlier EGs manually with gcloud dataplex entry-groups delete <name> --project <your-project> --location <your-location>.

The Acme Retail bundle is a synthetic OKF bundle for a US retailer's BigQuery estate. It contains nine leaf concepts across six directories (attesters, tables, metrics, computations, policies, skills), each with its own index.md, plus a bundle root with its own index.md and log.md. That's 17 pushed Entries in total; Dataplex auto-creates one <eg>_entry alongside, so gcloud dataplex entries list returns 18 rows.

After the push completes:

  • Every concept markdown file is a Knowledge Catalog Entry, discoverable by search across the whole project or organization, depending on IAM configuration.

  • The revenue-ytd Attested Computation appears in the console with its sanctioned SQL, its executor, its attester, its verification history, and the full concept body.

  • An analyst searching Knowledge Catalog for "revenue" finds Acme Retail's business definition alongside the BigQuery table it computes from, both under one permission model.

  • A downstream agent that already calls LookupContext for BigQuery table Entries retrieves the bundle's context by adding the OKF entry names to its resources list.

Further, metrics/revenue.md becomes an Entry with two Aspects. The full entries.get response (with view=ALL) looks like:

code_block
<ListValue: [StructValue([('code', '{\r\n "name": "projects/acme-analytics/locations/us-central1/entryGroups/acme-retail/entries/metrics/revenue",\r\n "entryType": "projects/acme-analytics/locations/us-central1/entryTypes/okf-bundle",\r\n "createTime": "2026-08-15T00:48:39.123456Z",\r\n "updateTime": "2026-08-15T00:48:57.234567Z",\r\n "parentEntry": "projects/acme-analytics/locations/us-central1/entryGroups/acme-retail/entries/metrics/index",\r\n "entrySource": {\r\n "displayName": "Revenue",\r\n "description": "Recognized revenue for a period, per Acme\'s FY2026 revenue-recognition policy. Backed by an Attested Computation.",\r\n "labels": {\r\n "finance": "true",\r\n "revenue": "true",\r\n "headline-metric": "true"\r\n },\r\n "location": "us-central1"\r\n },\r\n "aspects": {\r\n "dataplex-types.global.overview": {\r\n "aspectType": "projects/dataplex-types/locations/global/aspectTypes/overview",\r\n "createTime": "2026-08-15T00:48:57.111111Z",\r\n "updateTime": "2026-08-15T00:48:57.111111Z",\r\n "aspectSource": {},\r\n "data": {\r\n "content": "# Definition\\n\\nRevenue for a fiscal year is the sum of `net_amount` over orders that (a) reached `order_status = \'delivered\'`, (b) completed the 30-day return window, and (c) fall in the fiscal year by `order_ts`. Multi-currency orders are converted to USD at the `order_ts` daily reference rate. [^revenue-policy]\\n\\nThe sanctioned computation is [`computations/revenue-ytd.md`](../computations/revenue-ytd.md). Consumers MUST run and attest that computation rather than composing their own SUM. The attester rejects any receipt whose executed SQL does not match the sanctioned form.\\n\\n# Reporting cuts\\n\\n- **By fiscal year:** the sanctioned computation takes `year` as its sole parameter.\\n- **By channel or category:** these are approved narrations, not new metrics. Join the receipt\'s row-level result to `orders.channel` or to `order_lines` × `products.category` client-side. Do NOT rewrite the sanctioned SQL.\\n\\n# Trust and freshness\\n\\n- **Verified:** VP Finance sign-off on 2026-07-01, against the FY2026 policy.\\n- **Stale after 2026-12-31:** Finance re-issues the revenue recognition policy each January. Consumers of this concept after 2027-01-01 MUST re-verify the definition against the new policy before serving.\\n\\n[^revenue-policy]: Revenue Recognition Policy (FY2026)",\r\n "contentType": "MARKDOWN"\r\n }\r\n },\r\n "acme-analytics.us-central1.okf": {\r\n "aspectType": "projects/acme-analytics/locations/us-central1/aspectTypes/okf",\r\n "createTime": "2026-08-15T00:48:57.222222Z",\r\n "updateTime": "2026-08-15T00:48:57.222222Z",\r\n "aspectSource": {},\r\n "data": {\r\n "okf_type": "Metric",\r\n "generated": { "by": "reference_agent/gemini-2.5-pro", "at": "2026-06-30T14:00:00Z" },\r\n "verified": [ { "by": "human:jsmith@acme", "at": "2026-07-01T09:00:00Z" } ],\r\n "status": "stable",\r\n "stale_after": "2026-12-31T00:00:00Z",\r\n "sources": [\r\n {\r\n "id": "revenue-policy",\r\n "resource": "policies/revenue-recognition.md",\r\n "title": "Revenue Recognition Policy (FY2026)",\r\n "author": "human:jsmith@acme",\r\n "last_modified": "2026-06-15T00:00:00Z"\r\n }\r\n ]\r\n }\r\n }\r\n }\r\n}'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7f590bbb7990>)])]>

The overview Aspect holds the full body of revenue.md. The okf Aspect carries the structured signal fields, so agents get provenance, source, and OKF type in a form they can filter on directly instead of parsing markdown. Server-side searchEntries filters on the top-level scalar fields and on the scalar subfields of record fields; agents narrow further on the array-element subfields client-side after entries.get. (Aspects and EntryTypes are keyed by project number in real API responses and search predicates; the acme-analytics project ID is shown throughout for readability.)

What pushing your OKF to Knowledge Catalog enables

Once the bundle is in Knowledge Catalog, it provides two capabilities to any agent that reads from the catalog:

  • Discoverability across the organization. Agents find bundle concepts through the same searchEntries and LookupContext APIs they already use for cataloged data, so an OKF bundle appears alongside BigQuery tables and other resources in every query it matches.

  • Governance. Bundle Entries inherit IAM from the EntryGroup, so a single agent call returns exactly what the caller is permitted to read, with no parallel permission model to maintain.

Discoverability across the organization
OKF bundle Entries appear in searchEntries results alongside BigQuery tables and other cataloged resources, so an agent already querying the catalog picks up new bundles automatically. To retrieve a concept's body, trust signals, or linked concepts from a match, the agent moves to LookupContext and entries.get.

A LookupContext call looks like this:

code_block
<ListValue: [StructValue([('code', 'POST https://dataplex.googleapis.com/v1/projects/acme-analytics/locations/us-central1:lookupContext\r\n{\r\n "resources": [\r\n "projects/acme-analytics/locations/us-central1/entryGroups/acme-retail/entries/metrics/revenue"\r\n ],\r\n "options": { "format": "yaml", "context_budget": "8000" }\r\n}'), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7f590bbb7dd0>)])]>

The response is a single context field containing a pre-formatted YAML block. The block carries the entry's catalogEntry, its type, its description, its tags as labels, and its overview: the full markdown body of the concept, including its trust and freshness section. LookupContext does not render custom Aspects, so an agent that needs the structured OKF signal fields (okf_type, generated, sources, and the other ten) reads them with entries.get and view=ALL alongside the LookupContext call.

There is no repository clone, no manual Aspect merging, and no re-parse of frontmatter. The agent uses the same API call any Knowledge Catalog client already makes.

An agent traversing an OKF bundle typically follows a three-step flow. An agent that already knows the specific Entry names it needs skips step 1. An agent that already knows the target EntryGroup and wants to enumerate the bundle exhaustively substitutes entryGroups.entries.list for step 1.

  1. searchEntries returns candidate Entry names and descriptions. Its scope accepts a project or organization; narrowing within that scope happens through query terms, including aspect predicates like aspect:acme-analytics.us-central1.okf.okf_type=Metric.

  2. LookupContext on the top few Entry names (up to ten per call) returns the full concept body as pre-formatted YAML; context_budget caps the response size.

  3. entries.get with view=ALL on any Entry returns its structured OKF signals (okf_type, generated, sources, and the other ten) directly, which the agent can then filter or attest on.

When a concept's sources[] references another concept by path, the agent calls LookupContext on that Entry name to walk the reference.

The full response for the Revenue Entry:

code_block
<ListValue: [StructValue([('code', "resources:\r\n -\r\n catalogEntry: projects/acme-analytics/locations/us-central1/entryGroups/acme-retail/entries/metrics/revenue\r\n type: OKF Document\r\n description: Recognized revenue for a period, per Acme's FY2026 revenue-recognition\r\n policy. Backed by an Attested Computation.\r\n overview: |-\r\n # Definition\r\n\r\n Revenue for a fiscal year is the sum of `net_amount` over orders that (a) reached `order_status = 'delivered'`, (b) completed the 30-day return window, and (c) fall in the fiscal year by `order_ts`. Multi-currency orders are converted to USD at the `order_ts` daily reference rate. [^revenue-policy]\r\n\r\n The sanctioned computation is [`computations/revenue-ytd.md`](../computations/revenue-ytd.md). Consumers MUST run and attest that computation rather than composing their own SUM. The attester rejects any receipt whose executed SQL does not match the sanctioned form.\r\n\r\n # Reporting cuts\r\n\r\n - **By fiscal year:** the sanctioned computation takes `year` as its sole parameter.\r\n - **By channel or category:** these are approved narrations, not new metrics. Join the receipt's row-level result to `orders.channel` or to `order_lines` × `products.category` client-side. Do NOT rewrite the sanctioned SQL.\r\n\r\n # Trust and freshness\r\n\r\n - **Verified:** VP Finance sign-off on 2026-07-01, against the FY2026 policy.\r\n - **Stale after 2026-12-31:** Finance re-issues the revenue recognition policy each January. Consumers of this concept after 2027-01-01 MUST re-verify the definition against the new policy before serving.\r\n\r\n [^revenue-policy]: Revenue Recognition Policy (FY2026)\r\n labels:\r\n finance: 'true'\r\n revenue: 'true'\r\n headline-metric: 'true'"), ('language', ''), ('caption', <wagtail.rich_text.RichText object at 0x7f590bbb7e50>)])]>

Governance
Permissions on the EntryGroup use standard Knowledge Catalog IAM. An agent that names both a bundle concept and the BigQuery table it grounds against in one call receives both, each subject to its own existing access control list (ACL), so the response carries only what the caller is already permitted to read. There is no parallel permission model to maintain.

Reading agents use roles/dataplex.catalogViewer, which grants the read paths: entries.get, LookupContext, and searchEntries. The identity that runs kcmd push uses roles/dataplex.catalogEditor, which grants the write paths: entries.create and entries.patch. One EntryGroup per bundle-owning team is the multi-team pattern, and IAM on the EntryGroup cascades to its Entries.

LookupContext resolves the entry names it is given, up to ten per call, within a single location. It does not follow links out of a concept's body, so an agent that wants a referenced concept must name it explicitly. Place the bundle's EntryGroup in the same location as the data it describes to fetch both in one call.

Lifecycle

kcmd push is an idempotent upsert. Re-running is safe (no duplicates, no error), but every push writes every Entry. Concept deletes require an explicit kcmd delete on the Entry, or cleanup.ts to remove the whole EntryGroup at once; cleanup.ts deletes only the EntryGroup and its Entries, so the shared okf AspectType and okf-bundle EntryType stay in place for other bundles that reference them. For continuous ingestion in production, wire a CI job to kcmd push on every commit to the bundle repository, using a service-account credential with roles/dataplex.catalogEditor on the target EntryGroup.

Getting started

OKF defines what a trustworthy bundle looks like. Knowledge Catalog makes it reachable across the organization. To get started, check out the following resources:

  1. Read the OKF v0.2 spec and browse the Acme Retail bundle.

  2. Author a small bundle for one domain your team owns.

  3. Sync it into your Knowledge Catalog project using the sample code's setup.ts (which registers the resources) and push.ts (which delegates to kcmd).

  4. Point your existing agents at Knowledge Catalog. New context becomes reachable through the same LookupContext and searchEntries calls they already use.

Google Cloud Japan Blog

Malachyte がマネージド リアルタイム AI で小売業のコールド スタート問題を解決した方法

詳細を表示

※この投稿は米国時間 2026 年 8 月 11 日に、Google Cloud blog に投稿されたものの抄訳です。

未知の潜在顧客に商品をすすめるには、どのような方法が最適なのでしょうか。

私たちは、Spotify や Priceline などの大手企業とともに、長年この問題の解決に取り組んできました。そして、Sidd Motwani を中心に、AI を活用した e コマースのレコメンデーション プラットフォーム Malachyte を創設したのも、まさにこの問題を解決するためでした。最近では、消費者はパーソナライズされて関連性の高いコンテンツを期待するようになり、オンライン サービスは消費者の関心を集めるために競わねばならず、その結果がビジネスの成否を左右するまでになっています。

Malachyte は、高度な AI モデル、特に大規模言語モデルを、パーソナライズやレコメンデーションといった従来の課題に新しい方法で適用する独自のインサイトを見い出し、それが革新を生むきっかけとなりました。

さらに、潜在顧客の獲得を目指すうえで、Malachyte が常に思い描いていたビジョンを実現し、パーソナライズ アルゴリズムを構築し続けるには、セキュアかつスケーラブルで信頼性が高く、何よりも最先端の AI インフラストラクチャが必要でした。そこで、BigtableManaged Service for Apache Kafka などの Google Cloud ツールを活用することで、Malachyte は小売業者の売上を 2 倍、場合によっては 3 倍に増やすことに成功したのです。

ここでは、Malachyte がそのプラットフォームをいかにして構築したか、そして、企業家の皆様がこのようなサービスを利用して AI 基盤モデルのデプロイを新しい方法で開始するにはどうすればよいか、ご紹介します。

Malachyte はお客様の売り上げをどのように伸ばしたのか

Malachyte には、アテンション メカニズムを備えたニューラル ネットワーク(大規模言語モデルを支えるのと同じコンセプト)を使用して、小売検索と商品ページをパーソナライズできることを発見したときがひらめきの瞬間でした。このアプローチは、大規模言語モデル(LLM)でシーケンス内のアイテムの相対的な順序から意味を引き出すために使用されています。この場合、シーケンスは文中の単語と音節の順序です。Malachyte は、これを小売業のウェブサイトやアプリに応用し、顧客がサイトを操作する際の一連の行動の把握に努めました。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="1 - Malachyte blog" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/1_-_Malachyte_blog_.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>LLM が文の次の単語を予測するように、e コマースサイトでユーザーが次に何を求めているかを予測できたらどうでしょうか。</p></figcaption>
  
</figure>


  </div>
</div>

GPT 以前の言語モデルでも、単語や単語の断片(現在ではトークンとして知られているもの)の特定の並びを調べることは可能でした。しかし、それらの初期のモデルでは、単語が同等の順序で並んでいても、連続していなかったり、並べ替えられていたりした場合にどうなるかを調べることはできませんでした。この状況が打破されたのは、LLM がアイテムのシーケンス内の複雑で広い範囲の依存関係を理解できるようになったときでした。これには、Google の Transformer に関する取り組みも部分的に貢献しています。

このより洗練された手法は、生成 AI の一般的な普及はもちろん、Malachyte によるこのテクノロジーの応用という点でも、劇的な成果をもたらしました。

これを実践で機能させるために、Malachyte はユーザーがサイトにアクセスしたときに、そのユーザーについて知り得るすべての情報をベクトル化しています。ほとんどのユーザーは初めてサイトを訪れるため、ユーザーに関する情報はほぼありません。これが、いわゆる「コールド スタート」問題です。これに対応するには、ユーザーとのあらゆるやり取りをすべて活用して、このベクトルを改良する必要があります。たとえば、クリックやクエリなど、ベクトルに新しい要素が追加されるたびに、ユーザーが次に何を求めているかの予測が促進され、ユーザーに関する情報がさらに提供されます。

Malachyte のプラットフォームでは、その際にユーザー ベクトルと予測を同時に更新します。これにより、個々のユーザーとその好みをより深く理解できるだけでなく、匿名化されたユーザーデータを使用してモデル全体を改善することができます。推論を行うたびに、平均的な買い物客と特定の買い物客の両方のコンテキストが拡大する仕組みです。

Malachyte は、アテンションベースのニューラル ネットワークを使用するだけでなく、ユーザー プロファイルを 100 ミリ秒ごとに更新することで、さらなるイノベーションを実現しています。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="2 - Malachyte Blog" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/2_-_Malachyte_Blog.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>Malachyte のレコメンデーション エージェントと検索エージェントは、ユーザーが前のページでクリックした内容に基づいて、次のページの検索結果やレコメンデーション カルーセルにデータを入力します。</p></figcaption>
  
</figure>


  </div>
</div>

正確には、このコンセプト自体は新しいものではありません。小売業者は長年にわたり、コラボレーション フィルタリング レコメンデーション システムを使用して、類似するユーザーとアイテムを特定してきました。ただし、これには膨大な量のインタラクション履歴が必要でした。これらのモデルでは通常、サードパーティ Cookie ベースのプロファイルやユーザー属性など、大量のデータを必要とします。

しかし、セッション内のインタラクションのシーケンスに焦点を当てることで、小売業者はユーザーのプロファイルのみに焦点を当てるよりも、必要なデータや費用を抑えながら、はるかに多くのパーソナライズを実現できるようになります。また、長期的な Cookie データに依存しないことで、ユーザーのプライバシーをより尊重できる利点もあります。

これは、モデルの構造と、ブラウザデータ、クリック履歴、検索など、ユーザーに関するあらゆる情報をエンコードするマルチモーダル ベクトルによって実現されます。出力もマルチモーダルです。同じモデルを、サイト内検索の商品ページ、カテゴリページ、「カートに追加」カルーセルに適用できます。

この仕組みは、カタログ内の各商品がユーザー ベクトルと同じ空間に埋め込まれることで機能します。それにより、ユーザー ベクトルが更新されて、関連する商品に近づき、関連しない商品からは遠ざかるようになります。エンベディングを計算するニューラル ネットワークは、Malachyte と連携する小売業者全体で継続的にトレーニングされ、すべてのユーザーの品質を向上させています。このシステムは事実上、データ協同組合として機能し、各小売業者のユーザーがモデルをさらにスマート化させながら、それが全ユーザーに役立つ仕組みになっています。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="3 - Malachyte blog vector space - high res" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/3_-_Malachyte_blog_vector_space_-_high_res.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>商品の空間でベクトルとして表されるユーザー セッション。</p></figcaption>
  
</figure>


  </div>
</div>

この仕組みをすべてのユーザーに対して推論ごとに 100 ミリ秒で実現するために、Malachyte は Google Cloud のリアルタイム AI スタックを基盤とすることに大きなメリットを見出しました。

このシステムでは、すべての行動イベントが Managed Service for Apache Kafka クラスタにストリーミングされます。各イベントは、将来のトレーニング ジョブのためにキューに入れられるのではなく、Bigtable で直ちにユーザー プロファイルの更新に使用されます。Kafka クラスタにより、お客様のフロントエンドが永続化されるため、ユーザー セッションのシグナルをユーザー ベクトルにどのように適合させるかをあまり心配することなく、迅速に処理できます。

Bigtable によって、Malachyte のサービスは適切なユーザー ベクトルを検索して更新できるようになります。Bigtable と Kafka はステップごとに 10 ミリ秒のオーダーで動作するため、レコメンデーション ループ全体が、ユーザー エクスペリエンスを中断することなく完了します。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="4 - Malachyte Blog" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/4_-_Malachyte_Blog.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>小売業者のウェブサイト、Malachyte の AI モデルとサービング フロントエンド、コンテキスト管理インフラストラクチャから成る 3 層のリアルタイム AI アーキテクチャ。</p></figcaption>
  
</figure>


  </div>
</div>

高速コアの機能に加え、商品カタログの更新事項、在庫シグナル、小売店固有のディメンション データから成る第 2 レイヤには、商品データを最新の状態に保つ働きがあります。このデータは Cloud Pub/Sub 経由で送信されます。Cloud Pub/Sub は、グローバル アクセスが可能な REST API を提供するため、小売業者は緻密な統合作業なしで接続を利用できます。Malachyte エージェントは Google Kubernetes Engine (GKE)で実行され、モデル推論は Google Compute Engine(GCE)で行われます。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="5 Malachyte blog" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/5_Malachyte_blog.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>商品カタログの更新事項などの外部データの継続的な取り込みは、Pub/Sub のグローバル メッセージング システムを通じて行われます。</p></figcaption>
  
</figure>


  </div>
</div>

Google Cloud AI アーキテクチャに移行したことで、Malachyte は本番環境の AI 推論とトレーニングが GPU とストレージだけでは決まらないことを実証しました。高速な Key-Value ストア、ストリーミング レイヤ、マネージド メッセージング システムを含むリアルタイムの継続的学習インフラストラクチャが必要であり、これらすべてが基盤モデル アーキテクチャと統合されていることが重要です。

このアプローチは、Malachyte のような小規模なチームでも業界に多大な影響を与えることができることを示しています。必要なのは、強力なインフラストラクチャとコア AI マネージド サービスへのアクセスです。

使ってみる 

Malachyte のように業界に革新をもたらしたり、他社に先駆けて業界をリードしたりすることを目指すなら、Managed Service for Apache KafkaCloud Pub/SubBigtable をお試しください。新規のお客様は、$300 分の Google Cloud クレジットを利用できます。

- Malachyte、CEO、Sidd Motwani 氏

Malachyte、スタッフ ML エンジニア、Vicki Boykis 氏

Google Dataflow で費用対効果に優れた高スループットな生成 AI ワークフローを構築する

詳細を表示

※この投稿は米国時間 2026 年 8 月 19 日に、Google Cloud blog に投稿されたものの抄訳です。

リアルタイム ストリーミング パイプラインは現代の企業において運用のバックボーンとなっており、カスタマー サポートでのやり取りからトランザクション ログまで、あらゆるデータを処理し続けています。従来のストリーミング DAG は静的なものであり、一度デプロイされると、その処理ロジックや実行パスは変更できなくなります。しかし、生成 AI エージェントを統合すれば、静的なロジックから適応型の実行へと移行することが可能です。これにより、データの内容に応じてストリーミング ワークフローが実行時にプランを動的に立案し、データベースへのクエリや独自の修復処理を自律的に開始できるようになります。

たとえば、注文品の破損について顧客から不満を示すメッセージが届いた際、パイプラインの役割は単なるエラーログの記録や、ダッシュボードへのフラグ立てにとどまるべきではありません。顧客の注文と在庫の記録を保持するデータベースで注文情報を検索し、修復アクション(交換品の発送や払い戻しなど)を自ら判断して、顧客へのメール送信から最終的な対応結果の記録までを一貫して行うべきです。

しかし、生成 AI ワークフローをストリーミング システムに実行しようとすると、スケール、レイテンシ、コストという、エンジニアリング上の根本的な障壁に直面することになります。すべての未加工イベントを、外部データベースとメールツールを備えた高負荷なモデルやマルチステップ エージェントに直接送信すると、コストが跳ね上がり、レイテンシも増大します。また、API のレート制限もすぐに使い果たしてしまいます。

ここで紹介するパターンは、Google Dataflow(Google Cloud が提供する Apache Beam のフルマネージド サーバーレス実行サービス)と Agent Development Kit(ADK)を組み合わせてハイブリッド ストリーミング パイプラインを構築することで、スケーラビリティと複雑さという課題に対処します。軽量で CPU バウンドな ML モデルをアップストリームで使用してイベントをフィルタリングし、評価することで、パイプラインの費用対効果を高く保ち、複雑なケースのみをダウンストリーム エージェントにルーティングします。ケースを受け取ったエージェントは、実行するアクションを動的に決定し、パイプラインの静的 DAG に膨大な数の条件付きステップをハードコードすることなく、ストリームに動的分岐を導入します。

大規模ストリーム向けのユニバーサルなブループリント

以下ではカスタマー サポートにおけるトリアージのシナリオを例にしていますが、この「事前フィルタ + エージェントによるアクション」というパターンは、幅広い分野に適用できる普遍的モデルです。このパターンは、イベントの大多数(90% 超)が定型的であり、複雑なコンテキスト推論を必要とするイベントはごくわずかである、以下のようなストリームに適用できます。

  • IT 運用と DevOps: 数百万件の定型的なシステムログを CPU 上でフィルタリングし、重大な異常が検知された場合にのみエージェントをトリガーして診断を実行し、バグチケットを開く。

  • 金融詐欺のトリアージ: 数百万件のトランザクションを軽量なローカルルールでスクリーニングし、疑わしいパターンが検出された場合にのみエージェントを呼び出してマルチデータベース検索ツールを実行する。

  • 産業用 IoT: エッジで定常的なテレメトリーをモニタリングし、突発的な異常をエージェントにルーティングして、機器のシャットダウンやフィールド エンジニアへの通知など一連の対応を調整する。

アーキテクチャ: ストリーミング イベントを事前にフィルタする理由

高スループットなストリームの場合、メッセージの大部分は複雑な推論や修復を必要としません。その多くは、肯定的なフィードバック、中立的な問い合わせ、単純な質問などです。

すべてのイベントを高負荷な LLM ワークフローにルーティングすると、主に 3 つのボトルネックが発生します。

  1. API のコスト: フロンティア モデルはトークン単位で課金されます。高スループットの場合、コストはストリームのボリュームに比例して増加します。

  2. レイテンシ: データベース参照や外部 API 呼び出しを含むマルチステップ ワークフローには数秒かかるため、ストリーミング DAG のボトルネックになります。

  3. 割り当て: 外部 API には厳格なレート制限があり、ストリーミング ワーカーがその枠をすぐに使い切ってしまうおそれがあります。

これを防ぐため、Apache Beam と Dataflow を使用して、事前フィルタリングされたパイプラインを構築します。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="image1" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/image1_XYX8VCT.max-1000x1000.jpg" />
    
    </a>
  
</figure>


  </div>
</div>

パイプライン フロー

  1. 取り込み: Google Pub/Sub から未加工の顧客メッセージを読み取ります。

  2. 軽量の感情分類器(CPU): Apache Beam の RunInference 変換を使用して、軽量の CPU ベースの Hugging Face モデル(distilbert-base-uncased-finetuned-sst-2-english)ですべてのメッセージを処理します。これは Dataflow ワーカーの CPU でローカルに実行されるため、外部 API の費用は発生しません。

  3. 事前審査ゲート: シンプルな DoFn でストリームをフィルタします。感情が POSITIVE または NEUTRAL と判定されたメッセージは、確認応答を返したうえで破棄されます。

  4. 自動修復(ADK): メッセージが NEGATIVE に分類された場合にのみ、ADKAgentModelHandler を介して gemini-3.5-flash を基盤とする生成 AI エージェントをトリガーします。エージェントは各種ツールを活用して、BigQuery でのユーザー検索、注文情報の取得、対応策の選定、Gmail API を介した通知メールの送信までを自律的に実行します。

適応的実行: Beam DAG の動的な構成

従来のストリーミング アーキテクチャでは、パイプラインの有向非巡回グラフ(DAG)は固定されています。Dataflow にデプロイされると変換のシーケンスが固定されるため、新しいタイプのアラートを処理したり、特定のイベントのルーティング方法を変更したりする必要が生じた場合はパイプライン全体を変更、テスト、再デプロイする必要があります。

感情の事前フィルタのダウンストリームに生成 AI エージェントを配置することで、静的な DAG 内に動的で適応性のあるノードを導入できます。

全レコードの 95% を占める肯定的または中立的なデータは、パイプライン上の高速な静的パスに沿って処理されます。しかし、フィルタがネガティブなレコードを検知すると、エージェントはペイロードを評価し、実行時に API ツールの適切な実行順序(データベースへのクエリ、在庫確認、メール通知など)を動的に選択します。これにより、パイプラインで複雑な決定木を動的に実行できるようになり、静的な Apache Beam コードに膨大な数の条件分岐をハードコードして構築、保守し続けるといった煩雑な作業が不要になります。

パイプラインの実装

Google Agent Development Kit(ADK)と RunInference フレームワークを活用した Apache Beam での実装例を以下に示します。

1. 軽量な感情モデルの定義

HuggingFacePipelineModelHandler を使用して、アップストリームの CPU モデルを定義します。このモデルは、ワーカー インスタンス上で感情を POSITIVENEUTRALNEGATIVE に分類します。

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 0x7f0e934e2590>)])]>

2. 高負荷な ADK エージェントの構築

ADK エージェントは修復アシスタントとして機能します。このエージェントに、次の 3 つのツールを実装します。

  • lookup_user: BigQuery に顧客のメールアドレスをクエリします。

  • lookup_orders: BigQuery に顧客の注文と現在の商品在庫をクエリします。

  • send_email: 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 0x7f0e93db0f10>)])]>

LlmAgent を構成し、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 0x7f0e93db3590>)])]>

3. Dataflow DAG の組み立て

パイプライン全体を明確に定義します。アップストリームの感情推論の結果はフィルタリング ステップ(FilterNegativeADK)に直接渡され、フィルタリング ステップはダウンストリームの 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 0x7f0e932d1810>)])]>

コストとパフォーマンスのメリット

このフィルタリング ステップを導入すると、エンジニアリング面と運用面の両方で大きなメリットがあります。

1. 大幅なコスト削減

すべての受信イベントに対して料金を支払うのではなく、顧客のネガティブな感情が含まれる一部のメッセージ(通常は 5% 未満)に対してのみ Gemini の入出力トークンの料金を支払うだけで済みます。残りの 95% は CPU インスタンスでローカルに分類されるため、API の追加費用はかかりません。

2. 高いストリーミング スループット

Dataflow は、CPU による分類ワークロードを多数のインスタンスに分散して処理します。CPU 推論はミリ秒単位で完了するため、パイプラインを水平方向にスケールして高スループットのイベント ストリームを処理できます。ツール実行のためにリクエストごとに数秒を要する高負荷な LLM エージェントは、呼び出し回数を最小限に抑えることで、バックログの発生を防いでいます。

3. ネイティブな Apache Beam 統合

DAG にエージェントを追加するのに、複雑なオーケストレーション ロジックや手動のスレッドプールは一切必要ありません。ADKAgentModelHandler を Beam ネイティブの RunInference 変換と併用することで、並列ワーカー スレッド、バッチ処理、統合の処理が自動化され、コードベースをクリーンでメンテナンスしやすい状態に保てます。

重要ポイント

ストリーミング データは高速かつ大量であるのに対し、高負荷な生成 AI 推論は低速でコストがかかります。

Google Dataflow と ADK を使用して事前フィルタリングを組み込んだパイプラインを構築することで、ローカルの CPU ベースのモデルが持つ「低コストとスピード」、そして Gemini を活用したエージェントの「高度な自動化機能」という両方の長所を最大限に引き出すことができます。

完全なコードベースを参照してご自身でデプロイしてみる場合は、next-2026-demo GitHub リポジトリをご確認ください。


Apache Beam は Apache Software Foundation の商標です。

- グループ プロダクト マネージャー、Reza Rokni

- Google Cloud、ソフトウェア エンジニア、Danny McCormick

AI 時代の FinOps:エージェント向けの新しい柔軟な請求と費用管理

詳細を表示

※この投稿は米国時間 2026 年 8 月 27 日に、Google Cloud blog に投稿されたものの抄訳です。

AI がより複雑な業務を担うようになる中、ビジネスリーダーは、利益率と予算を守りながら、エージェントを活用して迅速なイノベーションを可能にするという新しい課題に直面しています。AI から投資対効果を得るには、技術の進化にあわせて 財務運用 (FinOps) とともに費用管理も進化させる必要があります。これにより、明確なコストの可視化、プロアクティブなコスト統制、ニーズに合った柔軟な支払いモデルの実現が可能になります。

本日、Gemini Enterprise 全体にわたるエージェント ワークロード、ならびに Google Antigravity in Gemini EnterpriseAndroid Studio といったデベロッパー ツール向けに、請求を柔軟にする支払いオプションと新しい費用管理ツールを発表します。

  • 柔軟な支払いオプション:既存の予測可能なユーザー単位のシート サブスクリプションと、Gemini Enterprise app の新しい従量課金オプションを組み合わせることができます。これにより、タスクの途中で制限に達することなく、エージェント ワークロードを実行できます。

  • デベロッパーへのアクセス権と AI の一括管理:Google Antigravity および Android Studio での AI 利用が Gemini Enterprise のサブスクリプションに含まれるようになりました(一部のお客様に提供中、他のお客様には順次提供予定)。これにより、管理負担を増やすことなくデベロッパーにより多くの機能を提供できます。Google Antigravity、プラットフォーム、アプリにわたる利用状況は、個別ライセンスや請求でサイロ化されず、単一のビューに集約されます。

  • 用量の増加に応じた割引:AI ワークロードが安定している、または増加している場合、 Flexible Savings Plans を利用することで、許容できる月額支出にコミットし、トークン費用を 10〜20% 削減できます。最低額や上限額の条件はなく、新しい請求サイロの管理も発生しません。

  • 統合された支出ガードレール:AI 支出やプロジェクトに月額の上限を設定し、エージェントの実行費用を見積もり、請求書に反映される前に予期しない予算の急増を検知できるようになりました。

Gemini Enterprise の支出をコントロールしながらチームに柔軟性を提供

企業ごとに組織運営の方法は異なります。また同じ企業内であっても、AI の利用方法はチームによって様々です。ビジネス ユーザーは日常的な生産性ツールを継続的に使用する傾向にある一方で、技術チームは AI エージェント ワークロードを短期間に集中して実行することがあります。

実際の業務の進め方に費用を適合させるため、以下の支払いおよびライセンスの選択肢を Gemini Enterprise の各機能と組み合わせることができます。

オプション

仕組み

費用最適化に役立つ理由

Gemini Enterprise app のユーザー単位シート サブスクリプション

ユーザーごとに一定の月額料金を支払います。これにはプロジェクト全体で共有される毎日の割り当てプールが含まれます。

予測可能な予算:日々安定した生産性ニーズを持つチーム向け。財務経理部に明確な月額費用を提供します。

【新機能】 従量課金制 Gemini Enterprise app 

*一部のお客様に提供中、他のお客様には順次提供予定

前払い契約や基本サブスクリプション料金はなく、標準モデル API レートに基づき、チームが消費したコンピューティングとトークンの分のみを厳密に支払います。

使用した分のみ支払い:支出は実際の利用状況に応じて自動的に増減するため、プロジェクトの需要が低下したときに利用されていない社員の料金を支払う必要がありません。

【Google Antigravity in Gemini Enterprise の 向け新機能】統合されたプール割り当て

日々の利用上限はプロジェクト全体でプールされ、ビジネスアプリ、デベロッパー ツール、カスタム エージェントはこの共有枠から消費します。プールされた割り当てが常に優先して消費されますが、管理者は超過利用を許可するかどうかを制御でき、許可した場合の超過分は従量課金レートで請求されます。

リソース利用の効率化:ビジネス ユーザーの未使用の制限枠を、デベロッパーやカスタム API エージェントの負荷の高いワークロードに自動的に、割り当てられるため、日々の制限枠が無駄になりません。

【今後提供開始】遅延実行価格

*一部のワークロード向けに提供予定

対象となるエージェント ワークロードを遅延実行として指定すると、Gemini Enterprise Agent Platform のインテリジェント スケジューラがオフピーク時間帯のキャパシティで実行します。

急がない業務の大幅な割引:AI ワークロードを分離されたオフピーク容量で実行できるため、推論コストを半額程度に抑えつつ、標準の割り当てを使わずにすむため、同じ予算内で大幅に多くのエージェント ワークロードを実行できます。

単一の Gemini Enterprise サブスクリプションで高度なエージェント ツールをデベロッパーに提供

Google Cloud は、Google Antigravity in Gemini Enterprise を提供開始しています。Google Antigravity in Gemini Enterpriseは高度なエージェント コーディングおよびエージェント構築機能を技術チームにもたらすエージェント ファーストのデベロッパー プラットフォームで、対象となるお客様の Gemini Enterprise サブスクリプションに含まれます。さらに、Android デベロッパーは、プロフェッショナルな Android 開発向けのエージェント IDE である Android Studio 内で、 Gemini Enterprise サブスクリプションに含まれる Google Antigravity の割り当てをネイティブに活用できます。

エージェント コーディングの費用効率を高めるため、各 Gemini Enterprise サブスクリプションに含まれるデベロッパー ツールの割り当てをプール化し、Google Cloud プロジェクト全体で利用できるようにしています。これにより、チームは購入済みのキャパシティを無駄なく活用できます。一元化されたガバナンスと管理を維持しながら、デベロッパーに高度なエージェント ツールを提供できます。Google Antigravity in Gemini Enterprise の新機能や、お客様による本番環境での活用事例の詳細については、こちら(英語)をご覧ください。 

Gemini Enterprise Flexible Savings Plans(FSPs)でよりスマートに予算を管理

企業の AI ワークロードが安定または増加している場合、Gemini Enterprise Flexible Savings Plans(FSPs)は、Gemini Enterprise 全体の利用に対してシンプルな利用額ベースのモデルを提供します。FSPs は、予算の柔軟性を維持しながらトークン費用を削減するように設計されています。

  • プログラムによる割引:Gemini Enterprise 全体の月額利用額に対し、1 年間の約定で 10% 割引、3 年間の約定で 20% 割引を受けられます。 

  • ペースに合わせて調整:最低利用額や上限額の条件がないため、現在のトラフィックに合った月額コミットメントを決め、利用量の増加に合わせて調整できます。 

  • Enterprise Agreement(EA)に対応: FSPs の利用額は、既存の Google Cloud EA の枠からシームレスに消化されるため、全体的なクラウドの約定枠を細分化することなく、各事業部門専用の予算管理が可能です。

Gemini Enterprise Flexible Savings Plans は、セルフサービスのお客様およびエンタープライズ契約をご利用のお客様に提供を開始しています。

財政規律を維持しながらチームに構築の自由を提供

ビジネスリーダーとしての目標は、AI の潜在的な価値を制限することではなく、管理されていない AI 費用によって発生する財務上および運用上のリスクを取り除くことです。エンジニアリング、マーケティング、運用の各チームにエージェントを活用したイノベーションの自由を与えつつ、それらのエージェントの挙動を信頼できる可視性と、予算を保護するセーフティネットを確保する必要があります。

このギャップを埋めるため、Google Cloud Billing Console 内に 3 つの明確な目標に基づいた堅牢なネイティブ ガバナンス ツールを構築しました。 

  1. プロジェクト開始前に費用試算Google Cloud Pricing Calculator を使用すると、ユーザー単位のライセンス、デベロッパー ツール、バックグラウンド エージェントのランタイム全体にわたる Gemini Enterprise の見込み費用を試算できます。プロジェクト業務の開始前に、必要な予算を確認できます。 
  2. 支出を細かく管理せずにしきい値を設定:プロジェクトチーム全体の日々の利用変動を管理するため、以下のツールにモニタリングを依頼することができます。 
    • 早期の異常検知:プロジェクトの AI 支出が通常より高くなる傾向を示した場合、システムは根本原因分析とともに異常値をフラグ付けし、増加の原因となっている上位 3 つの SKU を特定します。これにより、何が変化したのかを正確に把握できます。
<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="1 Jul22_Anomalies_Image1" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/1_Jul22_Anomalies_Image1.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>Billing Console に表示された初期異常アラートと、原因の主なSKUを強調表示した根本原因分析(RCA: the Root Cause Analysis)の内訳</p></figcaption>
  
</figure>


  </div>
</div>
  • プロジェクト レベルの支出上限:プロジェクトで明確な予算上の制限が必要な場合、Google Cloud Billing Console で直接、確実な月額支出上限を設定できます。プロジェクトが上限に達すると、エージェントの API 呼び出しが一時停止し、本番インフラストラクチャに影響を与えることなく予算を保護します。また、予算の 50%、80%、100% に達した際に送信される自動メールアラートにより、支出上限に対する進捗状況を把握できます。 
  • 超過利用コントロール:支出上限に達し処理が停止した場合でも、コンソールで 1 クリックで業務を再開できます。継続的な運用を優先する場合は、超過利用を有効にして超過利用分を従量課金レートへとスムーズに移行させることができます。この超過分は FSPs から直接引き落とされるため、超過単価には大幅な割引価格を維持できます。
<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="3 PAYG Overage Enabled" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/3_PAYG_Overage_Enabled.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>プロジェクトの従量課金による超過利用の有効化</p></figcaption>
  
</figure>


  </div>
</div>

3. ビジネス価値の可視化統合された請求レポートを FinOps エージェントと組み合わせて使用することで、予算の支出先に関するインサイトを自然言語で生成し、経営陣に ROI を簡単に提示できます。

<div class="article-module h-c-page">
  <div class="h-c-grid">


<figure class="article-image--large
  
  
    h-c-grid__col
    h-c-grid__col--6 h-c-grid__col--offset-3
    
    
  ">

  
  
    
    <img alt="Cost overview FinOps" src="https://storage.googleapis.com/gweb-cloudblog-publish/images/Billing_overage_-_dashboard_-_new_afternoo.max-1000x1000.png" />
    
    </a>
  
    <figcaption class="article-image__caption "><p>Google Cloud Console における AI 支出レポート</p></figcaption>
  
</figure>


  </div>
</div>

AI 費用の最適化をさらに深めるには

モデルとインフラストラクチャの費用、レイテンシ、パフォーマンスを最適化するフルスタックの FinOps 戦略を構築するには、詳細なアーキテクチャ仕様とフレームワークをご覧ください。

  • 動的な容量管理によりインフラストラクチャの制約を克服する方法: Google Kubernetes Engine および Google Compute Engine の機能を活用してコンピューティング投資を最適化する方法をご紹介します。これらの機能は、リソースを自動的にスケジュールおよび再割り当てし、中断、過剰なプロビジョニング、特定のハードウェア構成への過度な依存を回避します。 

  • エンタープライズのお客様向けに Google Antigravity を拡張:技術チームがエージェント ファーストのワークフローによってソフトウェアのデリバリーをどのように加速しているか、デベロッパー ツールに関する詳細解説ををご確認ください。 

  • スポーツカーから学ぶ AI 支出最適化:トークン数が多ければ常に優れた AI になるとは限りません。Gemini Enterprise Agent Platform プロダクト マネジメント ディレクターである Mike Clark との対談で、馬力と効率のバランスをとり、AI への投資からより高いリターンを得る方法をご覧ください。 

  • 利用スパイク時の保護:負荷の大きいワークロードは、それ以外の時間にアイドル状態となる高価な専用インフラストラクチャの料金を支払うことなく、ピーク時に急増させることができます。AI の利用が拡大しても、Gemini モデルは人工的なレート制限に達することなくオンデマンドで自動的にスケーリングし、1 分あたり最大 5,000 万トークンを処理できます。プロビジョンド スループットの詳細をご覧ください。