Gemini Structured Extraction → Firestore Knowledge Graph Pipeline
Goal
Build a deterministic pipeline that ingests multi-source raw artifacts, extracts structured entities/events/relationships via Gemini JSON schema outputs, and upserts the results into a Firestore-backed operational knowledge graph.
1) Canonical ingestion contract
Use one immutable envelope for all sources before extraction:
{
"artifactId": "evt_20260212_001",
"tenantId": "tenant_abc",
"source": {
"system": "crm|support|analytics|call-center|upload",
"channel": "email|chat|voice|web|mobile",
"receivedAt": "2026-02-12T13:45:21Z"
},
"subject": {
"customerId": "cust_123",
"sessionId": "sess_123",
"accountId": "acct_456"
},
"payload": {
"contentType": "text/plain|application/json",
"rawText": "...",
"blobUri": "gs://bucket/path"
},
"provenance": {
"traceId": "trace_...",
"ingestedBy": "pipeline-service",
"version": "v1"
}
}
Required ingestion rules
- Keep raw payloads immutable (write-once).
- Store processing lifecycle separately (
status,attempt,error,updatedAt). - Preserve
artifactId+traceIdfor end-to-end lineage.
2) Structured extraction schema (Gemini)
Define an explicit JSON schema and require the model to respond against it.
Suggested response schema
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "JourneyExtraction",
"type": "object",
"required": ["entities", "events", "relationships", "signals"],
"properties": {
"entities": {
"type": "array",
"items": {
"type": "object",
"required": ["id", "type", "name", "confidence"],
"properties": {
"id": { "type": "string" },
"type": {
"type": "string",
"enum": ["CUSTOMER", "ACCOUNT", "PRODUCT", "CHANNEL", "AGENT", "CAMPAIGN", "ORDER", "ISSUE"]
},
"name": { "type": "string" },
"confidence": { "type": "number", "minimum": 0, "maximum": 1 },
"evidence": {
"type": "array",
"items": {
"type": "object",
"properties": {
"text": { "type": "string" },
"start": { "type": "integer" },
"end": { "type": "integer" }
}
}
}
}
}
},
"events": {
"type": "array",
"items": {
"type": "object",
"required": ["id", "type", "timestamp"],
"properties": {
"id": { "type": "string" },
"type": { "type": "string" },
"timestamp": { "type": "string", "format": "date-time" },
"attributes": { "type": "object", "additionalProperties": true }
}
}
},
"relationships": {
"type": "array",
"items": {
"type": "object",
"required": ["from", "to", "type", "confidence"],
"properties": {
"from": { "type": "string" },
"to": { "type": "string" },
"type": { "type": "string" },
"confidence": { "type": "number", "minimum": 0, "maximum": 1 },
"evidence": { "type": "string" }
}
}
},
"signals": {
"type": "object",
"properties": {
"sentiment": { "type": "string", "enum": ["NEGATIVE", "NEUTRAL", "POSITIVE", "MIXED"] },
"intent": { "type": "string" },
"journeyStage": { "type": "string" }
},
"additionalProperties": false
}
},
"additionalProperties": false
}
Generation settings
- Temperature:
0.0-0.2 - Top-p: low/controlled
- Enforce schema response mode and reject invalid output.
3) Two-stage extraction design
Use two distinct calls to improve precision and reduce schema drift:
- Stage A: summary + entity/event discovery - Detect principal actors, touchpoints, and timeline anchors.
- Stage B: relationship pass - Extract typed edges only, with confidence + evidence.
Merge stage outputs under one validated JourneyExtraction payload.
4) Firestore graph model
Use three collections with deterministic IDs.
kg_nodes/{nodeId}
{
"id": "cust_123",
"type": "CUSTOMER",
"name": "Jane Doe",
"aliases": ["J. Doe"],
"confidence": 0.97,
"lastSeenAt": "2026-02-12T13:45:21Z",
"sources": ["support_transcript", "crm_note"],
"updatedAt": "<serverTimestamp>"
}
kg_edges/{edgeId}
edgeId = hash(from + type + to)
{
"from": "cust_123",
"to": "issue_991",
"type": "REPORTS",
"confidence": 0.91,
"artifactId": "evt_20260212_001",
"timestamp": "2026-02-12T13:45:21Z",
"updatedAt": "<serverTimestamp>"
}
kg_events/{eventId}
{
"id": "event_001",
"type": "SUPPORT_CONTACT",
"timestamp": "2026-02-12T13:44:10Z",
"subject": "cust_123",
"attributes": {
"channel": "chat",
"resolution": "unresolved"
},
"artifactId": "evt_20260212_001",
"updatedAt": "<serverTimestamp>"
}
5) Write path (idempotent)
- Validate model output against schema.
- Start Firestore transaction/batch.
- Upsert nodes (merge by stable IDs).
- Upsert edges by deterministic hash ID.
- Upsert events by event ID.
- Append provenance references (
artifactId,traceId, extractor version).
6) Recommended controls
- Confidence gates:
- auto-write >= 0.8
- review queue 0.5–0.79
- discard/log < 0.5
- Drift detection:
- track schema validation failure rate
- monitor relation-type cardinality spikes
- Auditability:
- persist evidence spans/text snippets per node/edge
7) Minimal implementation checklist
- [ ] Add ingestion envelope contract in pipeline service.
- [ ] Add Gemini response schema and strict validation.
- [ ] Implement two-stage extraction orchestration.
- [ ] Implement idempotent Firestore upsert layer for nodes/edges/events.
- [ ] Add confidence-based routing (auto-write vs review queue).
- [ ] Add metrics: extraction latency, validation failures, edge/node write counts.
- [ ] Add replay job for failed artifacts.
8) Rollout plan
- Phase 1 (shadow mode): run extraction + validation, no writes.
- Phase 2 (canary): write low-risk entity types in one tenant.
- Phase 3 (general): full graph write with monitoring + replay.
9) Notes for this repository
- Keep schema and graph write logic in separate modules to avoid tight coupling.
- Prefer deterministic IDs and merge semantics across all graph writes.
- Track extractor version so reprocessing can compare model generations over time.
10) Concrete Gemini request templates
Node.js (Gemini structured output)
import { GoogleGenAI, Type } from "@google/genai";
const client = new GoogleGenAI({ apiKey: process.env.GEMINI_API_KEY! });
const extractionSchema = {
type: Type.OBJECT,
required: ["entities", "events", "relationships", "signals"],
properties: {
entities: {
type: Type.ARRAY,
items: {
type: Type.OBJECT,
required: ["id", "type", "name", "confidence"],
properties: {
id: { type: Type.STRING },
type: { type: Type.STRING },
name: { type: Type.STRING },
confidence: { type: Type.NUMBER }
}
}
},
events: { type: Type.ARRAY, items: { type: Type.OBJECT } },
relationships: { type: Type.ARRAY, items: { type: Type.OBJECT } },
signals: { type: Type.OBJECT }
}
};
export async function extractJourney(rawText: string) {
const response = await client.models.generateContent({
model: "gemini-2.5-pro",
contents: [
{
role: "user",
parts: [{ text: `Extract customer journey graph fields from:\n${rawText}` }]
}
],
config: {
temperature: 0.1,
responseMimeType: "application/json",
responseSchema: extractionSchema
}
});
return JSON.parse(response.text ?? "{}");
}
Python (Pydantic validation after model output)
from pydantic import BaseModel, Field
from typing import List, Dict, Any
class Entity(BaseModel):
id: str
type: str
name: str
confidence: float = Field(ge=0, le=1)
class Event(BaseModel):
id: str
type: str
timestamp: str
attributes: Dict[str, Any] = {}
class Relationship(BaseModel):
from_: str = Field(alias="from")
to: str
type: str
confidence: float = Field(ge=0, le=1)
class JourneyExtraction(BaseModel):
entities: List[Entity]
events: List[Event]
relationships: List[Relationship]
signals: Dict[str, Any]
11) Firestore upsert pattern (transaction-safe)
import { createHash } from "node:crypto";
import { FieldValue, Firestore } from "@google-cloud/firestore";
const db = new Firestore();
const edgeId = (from: string, type: string, to: string) =>
createHash("sha256").update(`${from}|${type}|${to}`).digest("hex").slice(0, 32);
export async function writeExtraction(result: any, artifactId: string, traceId: string) {
const batch = db.batch();
for (const node of result.entities) {
const ref = db.collection("kg_nodes").doc(node.id);
batch.set(
ref,
{
...node,
updatedAt: FieldValue.serverTimestamp(),
lastArtifactId: artifactId,
lastTraceId: traceId,
},
{ merge: true }
);
}
for (const rel of result.relationships) {
const id = edgeId(rel.from, rel.type, rel.to);
const ref = db.collection("kg_edges").doc(id);
batch.set(
ref,
{
...rel,
artifactId,
traceId,
updatedAt: FieldValue.serverTimestamp(),
},
{ merge: true }
);
}
for (const event of result.events) {
const ref = db.collection("kg_events").doc(event.id);
batch.set(
ref,
{
...event,
artifactId,
traceId,
updatedAt: FieldValue.serverTimestamp(),
},
{ merge: true }
);
}
await batch.commit();
}
12) Suggested repository placement
functions/src/kg/schema/:- Gemini response schema definitions
- JSON schema validator wrappers
functions/src/kg/extract/:- Stage A and Stage B prompt builders
- Gemini client adapters
functions/src/kg/store/:- Firestore node/edge/event upsert logic
- Deterministic ID/hash helpers
functions/src/kg/jobs/:- Queue worker entrypoints
- Retry/replay handlers
13) Definition of done (DoD)
- Schema validation success >= 99% on shadow traffic for 7 consecutive days.
- Duplicate edge-write rate < 0.1% (idempotency check).
- P95 extraction+write latency within agreed SLO.
- Alerting in place for validation spikes, queue lag, and write errors.
- Backfill/replay runbook tested end-to-end in non-prod.