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

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

3) Two-stage extraction design

Use two distinct calls to improve precision and reduce schema drift:

  1. Stage A: summary + entity/event discovery - Detect principal actors, touchpoints, and timeline anchors.
  2. 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)

  1. Validate model output against schema.
  2. Start Firestore transaction/batch.
  3. Upsert nodes (merge by stable IDs).
  4. Upsert edges by deterministic hash ID.
  5. Upsert events by event ID.
  6. Append provenance references (artifactId, traceId, extractor version).

7) Minimal implementation checklist

8) Rollout plan

9) Notes for this repository

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

13) Definition of done (DoD)

← All docsView source on GitHub →