Design & deliver solutions on Cognite Data Fusion
Level 2 made you productive by hand. Level 3 makes you a solution builder: architecture decisions that survive five years, everything as code in Git deployed by CI/CD to dev → test → prod, custom extractors, robust orchestration, AI agents that answer operational questions with citations, and the governance to run it all in production.
- Level 2 complete — you have personally built the "Unit 2 feedwater mini-solution" (pipeline, transformation, model, matching, notebook, dashboards, groups).
- Comfortable Python (packages, virtual environments, classes), Git basics (branch, pull request), YAML, and reading an API reference.
- Admin (or near-admin) rights in a dev CDF project, and a Git host with CI (GitHub Actions, Azure DevOps or GitLab). The
publicdatasandbox is read-only, so this level needs a project of your own for most deliverables.
The Level 3 roadmap
Solution architecture & storage strategy
⏱ 3 h- The reference architecture and where each zone's responsibilities start and stop
- How to lay out projects and environments
- Which CDF store to use for which data — and the cost of choosing wrong
- An identifier strategy that survives source-system changes
- How to turn a business pain into a phased, valued roadmap
1.1 The reference architecture, zone by zone
| Zone | Owns | Typical decisions | Failure mode to design against |
|---|---|---|---|
| Plant / OT | Sources, extractor hosts, firewall rules | Where extractors run (DMZ VM vs. next to historian), buffer sizing, who holds service-principal secrets | Link outage → data gap. Buffer + state store + heartbeat alerting (Level 2 M1). |
| CDF | Stores, pipelines, models, tools, security | Storage per data type (1.3), model layering (M2), projects/environments (1.2), everything as code (M3) | Click-ops drift between dev and prod. Toolkit + CI/CD. |
| Customer cloud | Lake, ML, enterprise BI | Which way data flows (connector, Spark), what is the system of record for each dataset | Two copies of the truth. Decide: CDF is the contextualised operations layer; the lake is the enterprise analytics layer. |
| Consumers | Adoption, use cases | Which tool for which persona; what is an app vs. a Canvas vs. a dashboard | Built but unused. Change management (M8). |
1.2 Project & environment strategy
| Environment | Purpose | Data | Who deploys | Toolkit validation-type |
|---|---|---|---|---|
hnm-dev | Build and break things | Subset / synthetic; refreshed from prod monthly | Anyone in the team, from a branch | dev |
hnm-test | Integration tests, UAT with domain experts | Full copy or representative slice | CI on merge to main | dev |
hnm-prod | Operations | Live | CI on a release tag, with approval | prod (extra checks, no destructive drops) |
- One project per plant or one enterprise project? One enterprise project with spaces and data sets per site is the default: one schema, one access model, cross-site analytics for free. Separate projects only for legal/regulatory isolation or radically different owners.
- Promotion is code, not copying. Nothing moves dev → test → prod except through Git (M3). Data is re-extracted per environment, not exported/imported.
- Naming that sorts.
<org>-<env>for projects;<org>_<domain>/<org>_<domain>_datafor spaces;ds_<source>_<site>for data sets;ep_,tr_,fn_,wf_prefixes for pipeline objects.
1.3 Choosing the right store
| Data | Store | Why | Don't |
|---|---|---|---|
| Sensor signals, calculated KPIs, manual readings | Time series (datapoints) | Billions of points, server-side aggregates, native in Charts/Grafana/Power BI | …put sensor values in data-model properties or RAW rows |
| Alarms & events logs, batch/manufacturing logs, well logs, audit trails | Records & streams | Designed for hundreds of billions of rows per year; typed via containers; mutable or immutable streams for different lifecycles | …model every alarm as a graph node (millions of nodes with no relations is the wrong shape) |
| Source-system tables before cleaning | RAW | Schemaless landing zone; replayable; cheap | …query RAW from apps — it is a staging area, not a serving layer |
| Things and their relations: assets, equipment, work orders, people, documents-as-entities | Data model instances | The knowledge graph; typed views; GraphQL; what tools and agents traverse | …store high-frequency or high-volume rows here |
| P&IDs, manuals, photos, CAD, point clouds | Files, 3D, 360° | Binary storage with contextualization services on top | …keep documents only in SharePoint with a link — the parser needs the bytes |
| Map features, pipelines/lines, sites | Geospatial | Spatial queries (within, intersects) | …encode coordinates as strings in metadata |
A city works because a factory is not built in a residential street and a hospital is not put in a warehouse. CDF's stores are zones: each is optimised for one traffic pattern. The architect's job is zoning — put every dataset where its query pattern lives, then the roads (APIs) are fast and the bills are small.
1.4 Identifier strategy
External IDs are forever: every relation, every dashboard, every agent citation points at them. Rules that have held up in production:
- Derive from the source, never generate. Tag, functional location, work-order number, document number. If the source has no stable key, hash the stable columns.
- Prefix by source system for non-plant objects:
sap:wo:4711,pi:U2.BFP1A.DISCH.PRESS. Plant tags stay as the plant knows them (BFP-1A) because people search for them. - Renames become aliases, not new IDs. The CDM's
aliasesproperty exists for this; Search finds both. - Merges keep the survivor. When two records prove to be one pump, re-point relations to the survivor and mark the other retired (a status property) rather than deleting — history references it.
- Case and whitespace are normalised once, in the transformation (
upper(trim(tag))), never ad hoc downstream.
1.5 From business pain to a valued roadmap
Repeat for 4–6 use cases, then look for the shared foundation (usually: the asset hierarchy, the historian tags, work orders, P&IDs). Phase 1 builds that foundation plus the single highest-value use case; every later phase is cheaper because the foundation exists. That is the mechanism behind the Forrester numbers in Level 1.
Write it for a real 2-unit thermal plant you know (anonymise if needed). Five pages, these headings: (1) Use cases & value hypotheses (4–6, one table); (2) Zone diagram like Figure 1.1 with every source and consumer named; (3) Store decisions per dataset (the 1.3 table filled in); (4) Project/environment/naming standard; (5) Identifier rules; (6) Security model summary (roles from Level 2 M8); (7) Phased roadmap with what each phase delivers and measures. Keep it — Modules 2–8 build exactly this.
Advanced data modeling
⏱ 5 h- Schema design trade-offs: containers, indexes, constraints, polymorphism, versioning
- The three-layer model (core → enterprise → solution) Cognite recommends
- Populating at scale without corrupting the graph
- Records & streams, NEAT, and query design for performance
2.1 Schema design that lasts
| Decision | Guidance | Why it matters later |
|---|---|---|
| Container per concept vs. one wide container | One container per concept (Pump, Motor), plus the CDM mix-ins. Wide "everything" containers are a data-lake habit. | Views can combine containers; containers cannot be split without a migration. |
| Property types | Use typed properties (float64, int64, timestamp, boolean, enum, direct relation, timeseries, file) — not text for everything. | Filters, aggregation and agents work on types; a "text" pressure limit cannot be compared. |
| Indexes | B-tree index on properties you filter/sort on (tag, status, timestamps); inverted index for list search. | Un-indexed filters scan; queries slow down as the plant grows. |
| Uniqueness constraints | On business keys (serial number per manufacturer). | Stops duplicate equipment from the second extractor. |
| Polymorphism | Interfaces for families: RotatingEquipment ← Pump, Fan, Compressor. Query the interface to get all. | Tools and agents can work on "all rotating equipment" without knowing each type. |
| Edges vs. direct relations | Direct relation for ownership/one-to-many (equipment → asset). Edge when the relation itself has data or is many-to-many (pipe connects A↔B with diameter). | Edges are heavier to write and query; direct relations cannot hold properties. |
| Reverse relations | Declare them on the parent view (asset.timeSeries) so consumers traverse both ways. | Otherwise every consumer re-implements the join. |
| Versioning | Add properties = same version. Remove/rename/retype = new view version, and keep the old one until consumers migrate. Model version bumps when the set of views changes. | Consumers (Power BI, apps, agents) bind to a version; breaking silently is the #1 outage cause. |
2.2 The three-layer model
ISA-95 (enterprise → site → area → work centre → work unit) shapes the asset hierarchy levels in the enterprise model; ISO 14224 (equipment classes, failure modes) shapes equipment types and failure-mode enums. Put the standard's codes in CogniteAssetClass / CogniteEquipmentType instances, not in property names — then a KKS-coded plant and an ISA-coded plant share one model.
2.3 Populating at scale
- Batch and chunk. The instances API accepts batches (the SDK chunks for you); on the write side, stream in chunks of ~1,000 nodes and let the SDK retry 429/5xx.
- Nodes before edges; parents before children only when you rely on
path/rootmaintenance — otherwise direct relations may dangle temporarily and resolve later. - Idempotent upserts with
existingVersionwhen two writers could race (optimistic concurrency); otherwise plain replace. - Deletes are rare and explicit. Prefer a
status: retiredproperty and a nightly workflow that hard-deletes what has been retired for 90 days. - Debug with the inspect endpoint — it tells you which views a node is "seen" through and why a property is missing (wrong container, wrong space).
# Bulk, resilient population from a pandas frame (SDK 7/8)
from cognite.client.data_classes.data_modeling import NodeApply, NodeOrEdgeData, ViewId, DirectRelationReference
from cognite.client.exceptions import CogniteAPIError
pump_view = ViewId("hnm_enterprise", "Pump", "1")
def to_node(row):
return NodeApply(space="hnm_plantA_data", external_id=row.tag,
sources=[NodeOrEdgeData(source=pump_view, properties={
"name": row.tag, "description": row.description, "serialNumber": row.serial,
"asset": DirectRelationReference("hnm_plantA_data", row.func_loc),
"ratedFlow": float(row.rated_flow) if row.rated_flow == row.rated_flow else None})])
nodes = [to_node(r) for r in df.itertuples()]
for i in range(0, len(nodes), 1000):
chunk = nodes[i:i + 1000]
try:
res = client.data_modeling.instances.apply(nodes=chunk, auto_create_direct_relations=True)
print(f"{i:>7}: created {sum(n.was_modified for n in res.nodes)} changed")
except CogniteAPIError as e: # log and continue; a workflow task will re-run the failed chunk
print("chunk", i, "failed:", e.code, e.message)
2.4 Records & streams
Records are typed rows (schema from containers, access from spaces) stored in a stream — outside the graph, so they never become nodes. Streams define lifecycle, not schema: a mutable stream allows update/delete with low-latency consistent reads; an immutable stream forbids updates but ingests very large volumes. Intended use: alarms & events, manufacturing/batch logs, well logs, historic archives — "hundreds of billions of records per year" scale.
| Question | Graph instance (node) | Record |
|---|---|---|
| Has relations that tools must traverse? | Yes → node | No (a record may reference an asset by external ID, but the graph does not traverse into it) |
| Volume | Thousands–millions | Millions–billions |
| Lifecycle | Long-lived, updated | Append-mostly; retention by stream |
| Example | Pump BFP-1A, work order WO-4711 | Every DCS alarm on BFP-1A (thousands a day) |
Pattern: alarms land as records in an immutable stream; a nightly Function aggregates "alarm count per asset per day" into a time series; the agent joins records to the pump by external ID when asked "how many alarms last week?".
2.5 NEAT — modelling with domain experts in Excel
NEAT (pip install cognite-neat) lets you author or import a model as a spreadsheet (sheets for properties, classes, containers) or from RDF/OWL, validate it, and push it to CDF as spaces, containers, views and a data model. Use it when the vocabulary is negotiated with process engineers who will never write DML — they edit rows in Excel, you review the generated schema in Git.
2.6 Query design for performance
- Filter push-down: put every filter in the query (
hasData, equals, range on indexed properties), not in pandas afterwards. - Cursor pagination always; never "limit 10000 and hope".
- Traversal depth: fetch one hop per request pattern (asset → time series), not five; deep traversals belong in a Function that materialises a result.
- Sync for consumers: the instances
/syncendpoint gives a cursor that returns only changes since last call — a downstream system stays in step with a small, cheap poll. - Aggregate server-side (count/avg/histogram over instances) for dashboards instead of listing everything.
- GraphQL for apps, instances API for pipelines: GraphQL is ergonomic and nested; the REST query endpoint exposes cursors, sync and precise filters.
hnm_enterprise v1 + a solution model- Design the enterprise model on paper first (interfaces, types, relations, indexes, enums from ISO 14224). Review it with one process engineer using NEAT's Excel form.
- Implement it in DML (Level 2 style, on the CDM), with an ISA-95 hierarchy and edge properties on
Pump —feeds→ Boiler(line size). - Create an immutable stream
alarmswith a container for DCS alarms; load one day of alarms as records. - Create solution model
PumpHealth v1withHealthScoreimplementing your enterprisePumpcontext. - Populate 5,000 synthetic nodes with the chunked script; measure "all sensors under a unit" query latency < 300 ms; document the indexes that made it so.
Cognite Toolkit: CDF as code & CI/CD
⏱ 4 h- How the Toolkit represents every CDF resource as YAML in Git
- The build → dry-run → deploy loop and the repository layout
- A CI/CD pipeline that deploys dev on push and prod on tag with approval
- Detecting drift and packaging reusable modules
3.1 Concepts
The Cognite Toolkit (pip install cognite-toolkit, command cdf) turns a CDF project into a Git repository. Resources are YAML files grouped into modules; a config.<env>.yaml selects modules and supplies variables per environment; cdf build renders templates into a build folder; cdf deploy applies only what changed. The same repository deploys dev, test and prod — only the config file differs.
3.2 The workflow and the repository
pip install cognite-toolkit
cdf auth init # writes a .env with CDF_CLUSTER, CDF_PROJECT, IDP_CLIENT_ID/SECRET, IDP_TENANT_ID (dev)
cdf auth verify # checks the service principal has what the Toolkit needs
cdf repo init # .gitignore, cdf.toml and friends for a new repository
cdf modules init # scaffold: pick modules from Cognite's library, or start empty (also: modules add / list / upgrade)
cdf build --config-yaml config.dev.yaml # renders modules/ + the chosen config -> build/ (--env is deprecated)
cdf deploy --dry-run # shows create / update / unchanged / delete per resource
cdf deploy # applies only the differences
cdf clean --dry-run # (dev only) remove resources no longer in the repo; --drop-data also purges data
cdf modules pull ... # bring a module edited in the UI back into YAML (reverse drift)
# verified against cognite-toolkit 0.8.66 — run cdf --help for your version
hnm-cdf/
├─ cdf.toml
├─ config.dev.yaml config.test.yaml config.prod.yaml
├─ .github/workflows/deploy.yml
└─ modules/
├─ hnm_foundation/ # shared: identity, data sets, enterprise model, pipelines
│ ├─ auth/ unit2-engineers.Group.yaml · data-engineers.Group.yaml · extractor-cmms.Group.yaml
│ ├─ data_sets/ datasets.DataSet.yaml
│ ├─ data_models/ hnm_enterprise.Space.yaml · Pump.Container.yaml · Pump.View.yaml · hnm_enterprise.DataModel.yaml
│ ├─ raw/ cmms.Database.yaml
│ ├─ extraction_pipelines/ ep_cmms_equipment.ExtractionPipeline.yaml · ep_cmms_equipment.config.yaml
│ ├─ transformations/ tr_assets.Transformation.yaml · tr_assets.Transformation.sql · tr_assets.Schedule.yaml
│ ├─ functions/ fn_running_hours.Function.yaml · fn_running_hours/handler.py · requirements.txt · daily.Schedule.yaml
│ └─ workflows/ wf_nightly_unit2.Workflow.yaml · wf_nightly_unit2.WorkflowVersion.yaml
└─ hnm_pump_health/ # a solution module: its own model, functions, app
├─ data_models/ ...
└─ streamlit/ pump_health.Streamlit.yaml · pump_health/main.py
Three files from that tree, in the exact shapes the Toolkit expects (variables in {{ }} come from config.<env>.yaml):
# config.prod.yaml
environment:
name: prod
project: hnm-prod
validation-type: prod # stricter checks; no destructive drops
selected:
- modules/hnm_foundation
- modules/hnm_pump_health
variables:
location_name: plantA
unit2_engineers_source_id: 3f9c1c2e-.... # Entra ID group object id in the PROD tenant
schedule_nightly: "0 2 * * *"
# modules/hnm_foundation/auth/unit2-engineers.Group.yaml
name: unit2-engineers
sourceId: '{{unit2_engineers_source_id}}'
metadata: { origin: cognite-toolkit }
capabilities:
- timeSeriesAcl: { actions: [READ], scope: { datasetScope: { ids: [ds_pi_{{location_name}}] } } }
- assetsAcl: { actions: [READ], scope: { datasetScope: { ids: [ds_cmms_{{location_name}}] } } }
- dataModelsAcl: { actions: [READ], scope: { spaceIdScope: { spaceIds: [hnm_enterprise, cdf_cdm] } } }
- dataModelInstancesAcl: { actions: [READ, WRITE], scope: { spaceIdScope: { spaceIds: [hnm_{{location_name}}_data] } } }
- projectsAcl: { actions: [LIST], scope: { all: {} } }
# modules/hnm_foundation/transformations/tr_assets.Transformation.yaml (+ tr_assets.Transformation.sql beside it)
externalId: tr_assets_{{location_name}}
name: assets:{{location_name}}:cmms
dataSetExternalId: ds_cmms_{{location_name}}
destination:
type: nodes
view: { space: hnm_enterprise, externalId: Pump, version: '1' }
instanceSpace: hnm_{{location_name}}_data
conflictMode: upsert
ignoreNullFields: true
authentication: # the transformation runs as a service principal, never a person
clientId: ${TRANSFORMATION_CLIENT_ID}
clientSecret: ${TRANSFORMATION_CLIENT_SECRET}
tokenUri: ${IDP_TOKEN_URL}
cdfProjectName: {{cdf_project}}
scopes: ${IDP_SCOPES}
# tr_assets.Schedule.yaml: externalId: tr_assets_{{location_name}} interval: '{{schedule_nightly}}' isPaused: false
{{variable}} comes from config.<env>.yaml and is rendered at build time — use it for names, IDs, schedules. ${ENV_VAR} is read from the environment at deploy time — use it for secrets. Secrets never enter Git or the build folder.
3.3 CI/CD
The Toolkit docs ship setup guides for GitHub Actions, Azure DevOps and GitLab; the shape is always the same — dry-run on pull request, deploy on merge, prod on tag with approval. An illustrative GitHub Actions workflow (adapt from the official template for your Toolkit version):
# .github/workflows/deploy.yml (illustrative — start from the Toolkit's GitHub setup guide)
name: cdf-deploy
on:
pull_request: { branches: [main] }
push: { branches: [main], tags: ['v*'] }
jobs:
dry-run:
if: github.event_name == 'pull_request'
runs-on: ubuntu-latest
env:
CDF_CLUSTER: ${{ vars.CDF_CLUSTER }}
CDF_PROJECT: hnm-test
IDP_TENANT_ID: ${{ vars.IDP_TENANT_ID }}
IDP_CLIENT_ID: ${{ secrets.TEST_CLIENT_ID }}
IDP_CLIENT_SECRET: ${{ secrets.TEST_CLIENT_SECRET }}
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with: { python-version: '3.11' }
- run: pip install cognite-toolkit
- run: cdf build --config-yaml config.test.yaml
- run: cdf deploy --dry-run | tee dryrun.txt # reviewers read this in the PR
deploy-test:
if: github.ref == 'refs/heads/main'
runs-on: ubuntu-latest
environment: test
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with: { python-version: '3.11' }
- run: pip install cognite-toolkit
- run: cdf build --config-yaml config.test.yaml && cdf deploy
env: { CDF_CLUSTER: "${{ vars.CDF_CLUSTER }}", CDF_PROJECT: hnm-test, IDP_TENANT_ID: "${{ vars.IDP_TENANT_ID }}",
IDP_CLIENT_ID: "${{ secrets.TEST_CLIENT_ID }}", IDP_CLIENT_SECRET: "${{ secrets.TEST_CLIENT_SECRET }}" }
deploy-prod:
if: startsWith(github.ref, 'refs/tags/v')
runs-on: ubuntu-latest
environment: prod # GitHub "required reviewers" = the approval gate
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with: { python-version: '3.11' }
- run: pip install cognite-toolkit
- run: cdf build --config-yaml config.prod.yaml && cdf deploy --dry-run && cdf deploy
env: { CDF_CLUSTER: "${{ vars.CDF_CLUSTER }}", CDF_PROJECT: hnm-prod, IDP_TENANT_ID: "${{ vars.IDP_TENANT_ID }}",
IDP_CLIENT_ID: "${{ secrets.PROD_CLIENT_ID }}", IDP_CLIENT_SECRET: "${{ secrets.PROD_CLIENT_SECRET }}" }
- Separate service principals per environment, each with only the capabilities the Toolkit needs there; prod's secret exists only in the CI vault.
- Pull-request reviews read the dry-run. "3 to create, 12 to update, 0 to delete" is the sentence that prevents outages.
- Releases are tags; the tag is what you roll back to (re-deploy the previous tag).
3.4 Drift & testing
- Drift = someone changed prod in the UI. Detect it with a scheduled
cdf deploy --dry-runin CI that fails when anything would change; fix it bycdf modules pullinto a branch (if the change was right) or by re-deploying (if it was not). - Validation:
cdf buildvalidates YAML against the resource schemas; add a unit test that builds every environment on every PR. - Smoke tests after deploy: a small Python job (or a Function) that checks the transformation ran, the workflow succeeded, and a known GraphQL query returns rows.
3.5 Packaging modules
A module is reusable when everything site-specific is a variable: location name, data-set IDs, group source IDs, schedules, RAW database names. Test it by deploying the same module twice with two variable sets (plant A, plant B) into dev. A colleague should be able to onboard a plant by adding a config.<plant>.yaml and a few group IDs — nothing else.
cdf repo init+cdf modules initan empty repo;cdf modules pull/ hand-write every Level 2 object (groups, data sets, RAW, pipeline + config, transformation + SQL + schedule, function + schedule, workflow, spaces/containers/views/model) intomodules/hnm_foundation.- Three config files;
cdf build --config-yaml config.dev.yaml→cdf deploy --dry-run→cdf deployinto a fresh dev project. Everything from Level 2 must reappear. - CI: dry-run on PR, deploy on merge, prod on tag with a required reviewer.
- Break something in the UI on purpose; watch the scheduled drift check fail; fix via
cdf pull. - Deploy the module a second time with
location_name: plantB.
Building custom extractors
⏱ 3 h- The anatomy of a production extractor and the
cognite-extractor-utilspieces that give it to you - Incremental extraction, backfill, schema drift and dead letters
- When a hosted extractor beats writing code
- OT specifics and day-2 operations
4.1 Anatomy of a production extractor
# vendor_api_extractor/config.yaml
cognite:
host: https://westeurope-1.cognitedata.com
project: hnm-dev
idp-authentication:
token-url: https://login.microsoftonline.com/${TENANT_ID}/oauth2/v2.0/token
client-id: ${CLIENT_ID}
secret: ${CLIENT_SECRET}
scopes: [ https://westeurope-1.cognitedata.com/.default ]
extraction-pipeline: { external-id: ep_vendor_vibration }
extractor:
state-store: { raw: { database: extractor_state, table: vendor_vibration } } # state in CDF RAW: survives host rebuilds
upload-interval: 30 # seconds
upload-queue-size: 50000
source:
base-url: https://vendor.example.com/api/v2
api-key: ${VENDOR_API_KEY}
page-size: 500
metrics:
push-gateways: [ { host: http://prometheus-pushgateway:9091, job-name: vendor_vibration } ]
# vendor_api_extractor/extractor.py — a minimal but production-shaped extractor on cognite-extractor-utils
from dataclasses import dataclass, field
from datetime import datetime, timezone
import requests
from cognite.extractorutils import Extractor
from cognite.extractorutils.configtools import BaseConfig, StateStoreConfig, MetricsConfig
from cognite.extractorutils.uploader import RawUploadQueue, TimeSeriesUploadQueue
from cognite.client.data_classes import Row
@dataclass
class SourceConfig:
base_url: str
api_key: str
page_size: int = 500
@dataclass
class ExtractorSettings: # the "extractor:" section of config.yaml (hyphens map to underscores)
state_store: StateStoreConfig = field(default_factory=StateStoreConfig) # framework finds this and builds the store
upload_interval: int = 30
upload_queue_size: int = 50000
@dataclass
class VendorConfig(BaseConfig): # BaseConfig gives cognite:, logger:, version:, type: — the rest you declare
source: SourceConfig = None
extractor: ExtractorSettings = field(default_factory=ExtractorSettings)
metrics: MetricsConfig | None = None # "metrics:" section → Prometheus / Cognite pushers
def run_extractor(client, states, config: VendorConfig, stop_event):
"""Called by the framework with an authenticated client, the state store, parsed config and a stop signal."""
raw_q = RawUploadQueue(client, max_queue_size=config.extractor.upload_queue_size, max_upload_interval=config.extractor.upload_interval)
ts_q = TimeSeriesUploadQueue(client, max_queue_size=config.extractor.upload_queue_size, max_upload_interval=config.extractor.upload_interval, create_missing=False)
with raw_q, ts_q: # queues flush on exit
low, high = states.get_state("vendor:vibration") # watermark from the state store (None on first run)
since = datetime.fromtimestamp(high / 1000, tz=timezone.utc) if high else datetime(2024, 1, 1, tzinfo=timezone.utc)
url = f"{config.source.base_url}/measurements"
params = {"since": since.isoformat(), "limit": config.source.page_size}
while not stop_event.is_set():
page = requests.get(url, params=params, headers={"X-API-Key": config.source.api_key}, timeout=30).json()
for m in page["items"]:
ts_ms = int(datetime.fromisoformat(m["timestamp"]).timestamp() * 1000)
raw_q.add_to_upload_queue("vendor", "vibration_raw", Row(m["id"], m)) # keep the original in RAW
ts_q.add_to_upload_queue(external_id=f"vendor:{m['sensor']}:rms", datapoints=[(ts_ms, float(m["rms"]))])
states.expand_state("vendor:vibration", ts_ms, ts_ms) # advance the watermark
if not page.get("next"): break
params["cursor"] = page["next"]
with Extractor(name="vendor_vibration", description="Vendor vibration API → CDF", config_class=VendorConfig,
run_handle=run_extractor, version="1.0.0") as extractor:
extractor.run() # handles auth, config, state store, heartbeat to ep_vendor_vibration, metrics, graceful stop
Class and method names verified against cognite-extractor-utils 7.13 (which still requires cognite-sdk<8 — pin the SDK accordingly): Extractor(name, description, version, run_handle, config_class), RawUploadQueue.add_to_upload_queue(database, table, Row), TimeSeriesUploadQueue.add_to_upload_queue(external_id=…, datapoints=[(ms, value)]), state store get_state / expand_state. The framework evolves — re-check on upgrade.
4.2 Patterns that separate a script from an extractor
| Pattern | Implementation |
|---|---|
| Incremental with watermark | Persist the highest timestamp/ID per stream in the state store; query "since watermark − overlap" (overlap catches late arrivals); dedupe by key. |
| Backfill | A separate mode/flag that walks history in windows (day by day), writes the same queues, and updates the watermark only at the end — so a crash restarts the window, not the whole history. |
| Schema drift | Land the raw record untouched in RAW (as above); map only known fields to typed resources; alert when unknown fields appear rather than crashing. |
| Dead letters | Records that fail mapping go to <source>_dlq RAW table with the error; a daily check reports the count. |
| Idempotency | Deterministic external IDs (vendor:<sensor>:rms); RAW keys from the source ID; datapoints overwrite by timestamp. |
| Graceful stop | Honour the stop event; flush queues; save state — the framework's context manager does this if you loop on stop_event. |
4.3 Hosted extractors — when not to write code
If the source is an MQTT broker, Kafka topic, Azure Event Hub or a REST API that CDF can reach, configure a hosted extractor instead: a source (connection + credentials), a job (topic/endpoint + a mapping that turns incoming JSON into datapoints, RAW rows or instances), and a destination. No host, no upgrades, monitored like any pipeline. Choose code only when the source is behind the plant firewall, needs a proprietary client, or needs logic a mapping cannot express.
4.4 OT specifics
- OPC UA: prefer subscriptions (server pushes changes) over polling; set publishing interval and queue size per subscription; use the extractor's history-read for backfill; expect certificate-trust dances on every new server.
- PI: the PI extractor streams points; the PI AF extractor brings the AF hierarchy (assets) — run both, then entity matching is mostly unnecessary because AF already links tags to elements.
- Store-and-forward: enable the local buffer and monitor its size; a full buffer is an incident before it becomes data loss.
- Time and units: convert to UTC at the edge; carry the source unit into the time series'
sourceUnitand map to a CDM unit so Charts can convert.
4.5 Day-2 operations
- Containerise (a 15-line Dockerfile: python base,
pip install -r requirements.txt, entrypoint) and run as a systemd service or in the plant's container host; config and secrets injected as environment variables. - Heartbeat: the framework posts
seenruns; alert when none for 2× the expected interval. - Versioning: tag releases; the extractor reports its version to the pipeline; roll forward with a new container tag, roll back the same way.
- Extractor config in CDF: keep the YAML in the extraction pipeline's config (Toolkit
.config.yaml) so it is versioned with everything else and the extractor can fetch it at start-up.
- Stand up a fake vendor API (a 30-line FastAPI app returning paged vibration readings) and build the extractor above against it: RAW copy + time series, RAW state store, Prometheus metrics, registered pipeline with failure notifications.
- Containerise it; run it as a service; kill it mid-run and prove it resumes from the watermark without duplicates.
- Configure a hosted MQTT extractor against a public test broker (or a local Mosquitto) with a mapping to time series; compare effort.
- Put both pipelines' YAML into your Toolkit repo.
Orchestration & operations at scale
⏱ 3 h- Dynamic and sub-workflows, event triggers and versioning
- Running hundreds of Functions without surprises
- An observability stack for a CDF estate
- Data-quality gates and the cost/performance levers
5.1 Workflows, advanced
{ "workflowExternalId": "wf_nightly_all_units", "version": "v2",
"workflowDefinition": { "tasks": [
{ "externalId": "list", "type": "function",
"parameters": { "function": { "externalId": "fn_list_units" } } },
{ "externalId": "fanout", "type": "dynamic", "dependsOn": [{ "externalId": "list" }],
"parameters": { "dynamic": { "tasks": "${list.output.tasks}" } } }, // fn_list_units returns a list of task definitions,
{ "externalId": "report", "type": "function", "dependsOn": [{ "externalId": "fanout" }], // one subworkflow task per unit
"parameters": { "function": { "externalId": "fn_quality_report" } }, "onFailure": "skipTask" }
] } }
- Triggers: a cron schedule, or a data-modeling trigger (run when instances matching a filter change — e.g. new work orders arrived). Event-driven beats polling for "as soon as" requirements.
- Versions:
v1,v2live side by side; switch the trigger when v2 is proven; keep run history per version. - Inputs/outputs: pass small JSON between tasks; pass large data by reference (write to RAW/time series, pass the key).
5.2 Functions at scale
| Concern | Practice |
|---|---|
| Packaging | One folder per function with handler.py + requirements.txt, pinned versions; shared code as a small internal package installed from a wheel in requirements.txt. |
| Secrets | Function secrets (Toolkit: ${...} at deploy) — never requirements.txt URLs with tokens, never hard-coded. |
| Limits | Memory and run-time limits exist and differ per cloud (see the docs' technical page). Design for < 10 minutes; split by unit/site and fan out from a workflow. |
| Cold starts | Keep dependencies small; avoid importing heavy ML stacks for simple jobs. |
| Idempotency | A re-run must produce the same result: upsert, deterministic IDs, write windows by time range. |
| Logs & calls | Every call has an ID, logs and status; return a small dict summarising work done (rows, ranges) so the workflow run shows it. |
5.3 Observability
| Signal | Source | Alert when | Where it shows |
|---|---|---|---|
| Data freshness | Latest datapoint timestamp per data set (a Function computes it hourly into a time series) | > 2× expected interval | Grafana ops dashboard · Charts alert |
| Pipeline health | Extraction-pipeline runs (seen/success/failure) | failure, or no seen for N | CDF notifications (email/Teams) |
| Transformation results | Run history: rows written, errors | errors > 0, rows = 0 unexpectedly | Workflow run view; ops dashboard |
| Workflow duration | Per-task timings | > baseline × 1.5 | Ops dashboard |
| Extractor metrics | Prometheus pusher (queue size, rows/s, buffer bytes) | buffer growing, throughput 0 | Grafana |
| Quality | Quarantine counts, coverage % | quarantine spike; coverage drop | Weekly report + dashboard |
A runbook per alert (what it means, first checks, who to call, how to re-run safely) turns an alert into a 10-minute fix instead of an evening. Store runbooks next to the module in Git.
5.4 Data-quality gates
- Validate in the transformation: rows with null keys, impossible values or unknown parents are routed to
<source>_quarantineRAW tables (a second transformation with the inverse filter). - Coverage KPIs as time series: % time series linked to an asset, % equipment with a serial number, % drawings parsed, % work orders resolved to an asset.
- Contract tests in CI: a GraphQL query per critical view with expected minimum counts, run after every deploy.
- Ownership: every data set names an owner who receives the quality report.
5.5 Cost & performance levers
- Aggregates, not raw datapoints, for anything longer than a day (Level 2 rule) — the biggest single lever for read volume.
- Rate limits: the SDK retries 429s with back-off; run heavy jobs off-peak; fan out sensibly (10 parallel functions, not 500).
- Batch sizes: datapoints in thousands per request, instances in ~1,000, RAW rows in thousands; too small = many calls, too large = timeouts.
- Granularity of stored KPIs: a 1-minute calculated KPI costs 60× a 1-hour one; store what decisions need.
- Capacity planning: tags × frequency × retention → datapoints/year; know the number before the historian backfill starts.
- Refactor
wf_nightly_unit2intowf_nightly_all_units v2with a dynamic fan-out and one sub-workflow per unit; make one unit's transformation flaky on purpose (retried) and one fatal (abort only that sub-workflow). - Add a quality task that quarantines bad rows and writes coverage KPIs as time series.
- Build the Grafana ops dashboard: freshness per data set, pipeline status, workflow durations, quarantine counts; two alert rules with runbooks in the repo.
- Write the capacity plan for the full plant (tags, frequencies, retention) as a one-page table.
Atlas AI: building industrial agents
⏱ 4 h- What an agent is made of — instructions, model, tools, skills — and why grounding on the graph works
- Building one in Agent Builder and as code with the SDK
- Evaluating and tuning with test cases
- Calling agents from your own apps, and the safety rules
6.1 Concepts
| Good agent use case | Poor agent use case |
|---|---|
| Troubleshooting: "what changed on X, what should I check first?" — many sources, needs synthesis, human decides | Anything a dashboard answers with one number ("current discharge pressure") |
| Document Q&A across manuals, datasheets, procedures | High-frequency control decisions (latency, safety) |
| Shift handover / incident summaries with citations | Questions whose data is not yet in CDF — the agent cannot invent it |
| Guided data exploration for non-experts | Autonomous write actions without confirmation |
Built-in tools (from the docs' tool library): Query knowledge graph and the newer Query tool (list/search/aggregate on views); Query time series data points; Answer document questions, Summarize documents, Analyze time series (preview), Examine data semantically (preview); integration: Call REST API, Run Python code, Call Function, Sandbox (private preview). Skills hold workflow-specific knowledge — vocabulary mappings ("BFP" = boiler feed pump), which views to use, which tools in which order — and load only when relevant.
6.2 Agent Builder
- Scope & instructions: identity, what it handles, universal rules for ambiguous data, output format, what needs confirmation. Keep workflow knowledge out of instructions and in skills.
- Model: start with a high-capability model to set a quality baseline; test smaller/cheaper ones against your evaluation set later.
- Tools: Query knowledge graph bound to your enterprise model and instance spaces; time series tool; document tools; a Function tool for a calculation only your code can do.
- Skills: one per workflow (troubleshooting, handover summary), reusable across agents.
- Test in the chat panel, then publish (label
published) so it appears in Canvas/Charts/Search.
6.3 Evaluation
Before tuning anything, write 10–20 representative questions with expected answers (what facts, which citations, what format). Run the evaluation after every change to instructions, skills, tools or model; a change that raises one score and drops another is not an improvement. Typical rubric: correctness of facts, presence of citations, followed the format, asked for confirmation when it should, latency, cost per answer.
6.4 Agents as code
# pip install "cognite-sdk>=8" # client.agents (upsert/chat) exists from SDK 7.9x; the runtime_version field needs SDK 8.x
from cognite.client.data_classes.agents import (
AgentUpsert, Message,
QueryKnowledgeGraphAgentTool, QueryKnowledgeGraphAgentToolConfiguration, DataModelInfo, InstanceSpaces,
QueryTimeSeriesDatapointsAgentTool, AskDocumentAgentTool, SummarizeDocumentAgentTool,
)
agent = AgentUpsert(
external_id="pump_troubleshooter",
name="Pump troubleshooter",
description="Explains what changed on a pump and what to check first, with citations.",
instructions=(
"You are a reliability engineer's assistant for HNM power plants. Scope: rotating equipment. "
"Always: identify the equipment node first, then pull the last 7 days of key signals, then work orders, then documents. "
"Never guess values; if data is missing say so. Format: Findings · Evidence (with sources) · Suggested checks. "
"Ask for confirmation before calling any tool that writes or triggers actions."),
model="azure/gpt-4o",
runtime_version="1.3.0",
labels=["published"],
tools=[
QueryKnowledgeGraphAgentTool(name="graph", description="Find equipment, sensors, work orders and documents",
configuration=QueryKnowledgeGraphAgentToolConfiguration(
data_models=[DataModelInfo(space="hnm_enterprise", external_id="Enterprise", version="1")],
instance_spaces=InstanceSpaces(type="all"))),
QueryTimeSeriesDatapointsAgentTool(name="signals", description="Read sensor datapoints for equipment"),
AskDocumentAgentTool(name="ask_doc", description="Answer questions from datasheets, manuals and procedures"),
SummarizeDocumentAgentTool(name="summarize", description="Summarize a document"),
],
)
client.agents.upsert(agent)
response = client.agents.chat(agent_external_id="pump_troubleshooter",
messages=[Message(content="What changed on BFP-1A in the last 7 days and what should I check first?")])
print(response.text) # the answer with citations
print(response.action_calls) # any tool calls that need confirmation before they run
Agents, skills and evaluations can be kept as files and deployed by the Cognite CLI in CI — the same dev → test → prod promotion as everything else. Version the instructions and the evaluation set together; a prompt change without an evaluation run is an untested deploy.
6.5 Using agents from your own apps
The chat call is the integration point: a Streamlit or Flows app (M8) calls the agent with the user's identity, shows the answer and citations, and — for action_calls — renders a "Confirm" button that only then executes the tool. Agents published to Canvas/Charts appear in those tools automatically; a Function can call an agent to produce a nightly summary that is written back as a document.
6.6 Safety & governance
- Permissions: an agent sees what the calling user may see; scope the user, not the agent.
- Read by default; any writing/acting tool sits behind confirmation.
- Citations are mandatory in instructions — an answer without evidence is rejected by the evaluation.
- Audit: log questions, tool calls and answers (the platform records runs); review weekly for wrong answers and add them to the evaluation set.
- Model choice is a governance decision (data residency, cost); Atlas AI's LLM benchmarking exists to make it evidence-based.
- Build the agent above (Builder first, then as code) against your enterprise model with three tools and one "troubleshooting" skill.
- Write 15 evaluation questions with expected findings and citations; run; iterate instructions until ≥ 80% pass.
- Publish it; use it from Canvas on a real trip; screenshot the cited answer.
- Try a smaller model against the same evaluation; record quality/latency/cost.
- Put agent, skill and evaluation files in the repo; deploy to test via CI.
Simulators, robotics & 3D pipelines
⏱ 3 h- How process simulators plug into CDF and run from data
- How robots become data sources and mission executors
- Running 3D as a pipeline, not a one-off upload
- Where geospatial fits
7.1 Simulator integrations
The routine's data sampling and validation step is what makes simulation-from-live-data trustworthy: it takes a window of datapoints, checks steady state (low variance) and validity ranges, and only then feeds the simulator. Results land as ordinary time series, so the "digital twin" is just more tags next to the real ones.
The same pattern applies to thermodynamic models (heat-balance, condenser, turbine cycle) with a custom connector: sample boiler/turbine tags → run the model → write expected heat rate, expected condenser vacuum → alert on the delta. The connector toolkit exists so the simulator need not be one Cognite already supports.
7.2 Robotics — InRobot
- Missions are planned on the 3D model: waypoints, actions (photo, thermal image, gauge read, gas sniff) at assets.
- Robot data sources land images, readings and scans in CDF linked to the asset at each waypoint — the same graph InField observations use.
- Automation: schedule missions like rounds; trigger a mission from an alert (bearing temp high → go look); feed results to an agent for anomaly triage.
- Prerequisites: a contextualised 3D model (7.3) and robot connectivity; safety rules from the site define where robots may go.
7.3 3D as a pipeline
- Model & revision: create a 3D model, upload a CAD/point-cloud file as a revision (SDK:
client.three_d.models.create→client.files.upload→client.three_d.revisions.create, or the Toolkit's 3D module); CDF processes it for streaming. - Contextualise automatically: a Function walks the CAD node tree, extracts tags from node names with the plant's regex, matches to assets, and writes
Cognite3DObjectlinks; a domain expert reviews the low-confidence tail. - Re-run per revision: the pipeline is a workflow: new revision published → contextualisation function → coverage KPI.
- 360° & point clouds: site walk-throughs as image collections with station links; point clouds for as-built comparisons.
7.4 Geospatial
For linear or distributed assets (grid lines, pipelines, wind farms, water networks) use geospatial feature types: features carry geometry and properties, support spatial queries (within, intersects, nearest), and link to assets. Substations and towers as points, lines as linestrings; combine with time series for "which towers are within 5 km of the fault?".
- Install DWSIM (free) and its Cognite connector on a Windows VM; upload a simple pump/pipe model; define a routine that samples 24 h of suction pressure and flow and returns expected discharge pressure.
- Run it from a workflow task nightly; write the expected series; add a Charts alert on actual − expected > 5%.
- Upload a small CAD model (any sample plant) and contextualise 20 nodes with a Function; show "Show in 3D" from Search.
- Optional: model one line of a grid as geospatial features and query "towers within 2 km".
Cognite Flows, governance & the Level 3 checkpoint
⏱ 3 h- What Cognite Flows is and how an app goes from idea to certified production
- Choosing between Streamlit, Flows and a custom front-end
- A governance operating model that keeps the estate healthy after year one
- Running many plants from one codebase, and getting people to use it
8.1 Cognite Flows — production-grade apps on the graph
Cognite Flows is Cognite's platform for building custom web applications that run inside CDF. It is deliberately AI-native: you build with an agentic coding tool (Claude Code, Cursor, GitHub Copilot…) from templates that already contain CDF authentication, data access and SDK utilities, plus "skills" the coding agent uses, and a checked-in spec that both you and the agent work from. Apps use modern web tooling (Vite), hot reload and automated tests, run locally against CDF, and are deployed interactively or by CI/CD with version and access control. An optional OAuth relay handles Entra ID sign-in.
The Flows Foundation modules are named for it: a design lane (human-centred: who uses it, what decision it supports, accessibility, no dead ends) and an engineering lane (secure: least-privilege data access, no secrets in the client, tested, observable, versioned, documented handover). "Vibe-coded" apps fail on the engineering lane — which is the point.
8.2 Which app technology?
| Streamlit in CDF | Cognite Flows | Custom front-end (JS SDK + React) | |
|---|---|---|---|
| Best for | Internal tools, analyses, prototypes in hours | Production apps for operators/engineers inside CDF | Portals, mobile, embedding CDF in an existing product |
| Skills | Python | Web + agentic coding; Python or TypeScript for data access | Full web engineering |
| Auth & hosting | Built in | Built in; OAuth relay optional | You own it |
| Governance | Light | Certification lifecycle | Your SDLC |
| Rule of thumb | Prototype in Streamlit → if users keep asking for it, rebuild as a certified Flows app → go custom only when it must live outside CDF. | ||
8.3 Governance operating model
| Artefact | Owner | Cadence |
|---|---|---|
| Standards doc: naming, identifier rules, modeling conventions, security roles (your Deliverable 1, kept current) | Data-platform lead | Reviewed quarterly |
| Source onboarding checklist: owner, data set, extractor choice, pipeline + alerting, quality rules, access, documentation | Integration engineer per source | Per source |
| Model change management: proposal → impact on consumers → version plan → PR → migration → deprecation date | Model steward + consumers | Per change |
Access review: admins, all-scoped capabilities, service principals | Security + platform lead | Monthly / quarterly |
| Quality & freshness report (M5) | Data-set owners | Weekly |
| Runbooks & on-call | Ops | Per alert |
| Release notes review (CDF "What's new") | Platform lead | Monthly |
8.4 Multi-plant operations
- One enterprise model, one Toolkit repo, N config files — a plant is a set of variables and group IDs (M3.5).
- Environment parity: dev/test/prod built from the same tag; a plant is onboarded in dev first, then promoted.
- Shared modules, local extensions:
hnm_foundationeverywhere;plantA_extrasonly where needed. - Upgrade cadence: Toolkit and SDK pinned per release; upgrade on a schedule, in dev first, with the drift check green. Note the current pins:
cognite-toolkit 0.8.xandcognite-extractor-utils 7.xstill requirecognite-sdk<8. - Vendor roadmap: with Schneider Electric's acquisition of Cognite (announced June 2026) and the planned integration into AVEVA CONNECT, review Cognite's "What's new" and AVEVA announcements monthly for naming, packaging and API changes.
8.5 Change management & adoption
- Start with the persona who feels the pain (the reliability engineer, not "the plant"); ship the Canvas/agent that answers their Monday question.
- Champions per unit, trained on Level 1 Module 5 and the Domain Expert Basics path; they run the weekly "what did you find" session.
- Measure usage: active users per tool, canvases created, agent questions per week, and the business KPI from Deliverable 1 — report both together.
- Retire the spreadsheet only after the CDF path has been the source of truth for one full cycle (a month-end, an outage).
- Prototype "Pump Health" in Streamlit (pump picker from the enterprise model, 30-day trends with work orders, health score from the solution model, "Ask the troubleshooter" via the agent).
- Rebuild it as a Flows app from a template with an agentic coding tool; check in the spec; add tests; run the four certification steps as far as your project allows.
- Deploy via CI to test; write the handover page (purpose, users, data dependencies, runbook, owner).
- Write the one-page governance standard and the source-onboarding checklist for your organisation.
8.6 Level 3 checkpoint quiz
1. Thousands of DCS alarms per day per asset belong in…
2. In the three-layer model, ISO 14224 equipment classes live in…
3. In the Toolkit, {{variable}} vs ${ENV_VAR}:
4. The CI stage a reviewer must read before approving a merge is…
5. An extractor is restarted after a crash. It avoids re-sending history because of…
6. Running the nightly job for 12 units with isolated failures is best done with…
7. Why can an Atlas AI agent cite its sources?
8. Before tuning an agent's instructions you should…
9. A simulator routine's sampling & validation step exists to…
10. Cognite Flows' application certification checks…
★ Level 3 portfolio — what you can now show a client
Each deliverable is a portfolio piece. Tick what exists:
- Solution architecture document for a 2-unit plant: zones, stores, environments, identifiers, security, valued roadmap (M1)
-
hnm_enterprise v1on the CDM with a solution model, an alarm stream, 5,000 nodes and measured query latency (M2) - A Git repository that deploys the whole estate to a fresh project; CI with dry-run on PR, deploy on merge, prod on tag; drift check; a second plant onboarded by config (M3)
- A containerised extractor on extractor-utils with resumable state, metrics and alerting, plus a hosted MQTT pipeline (M4)
- A fan-out workflow with quality gates, an ops dashboard with runbooks, and a capacity plan (M5)
- A published troubleshooting agent with an evaluation set ≥ 80% and a model comparison (M6)
- Expected-vs-actual from a simulator routine, a contextualised 3D model (M7)
- A Pump Health app (Streamlit → Flows) with handover docs, a governance standard and an onboarding checklist (M8)
You can walk into a plant, run a two-day discovery, produce the architecture document, stand up a governed dev/test/prod estate from a Git repository in a week, integrate the historian, SAP and drawings, and hand operators a Canvas, alerts and a cited troubleshooting agent — then teach the client's team to run it. Next: the 4-day CDF Practitioner bootcamp to do it under instruction with Cognite engineers, the Flows Builder certification, and the two advanced data-modeling exams.
Certifications for this level
Glossary additions for Level 3
dev/prod) that tightens checks and prevents destructive operations in production.