From OPC UA to Digital Twin: Time-Series Ingestion Architecture

From OPC UA to Digital Twin: Time-Series Ingestion Architecture

Digital Twin Data Ingestion: From OPC UA to a Time-Series Twin

Most digital twin projects do not fail in the modeling tool. They fail in the plumbing, usually around month four, when someone notices that the twin says a pump is running while the pump has been dark for forty minutes, or that last Tuesday’s vibration history has a hole exactly where the plant network was patched. The model was fine. The digital twin data ingestion path underneath it quietly lost, reordered, or duplicated the facts the model depends on.

This matters more in 2026 than it did when twins were dashboards. A twin that feeds a simulation, a predictive-maintenance model, or an agent that proposes setpoint changes inherits every ingestion defect as a wrong decision, not a wrong chart. The question stops being “can we get the data” and becomes “can we prove what the data means, when it was true, and what we know we are missing.”

By the end of this article you will have a reference architecture for OPC UA to digital twin ingestion, concrete client and database code, a way to reason about store-and-forward, late data and idempotency, and a split between the history of a signal and the state of a twin that most designs blur.

What this covers: the sources and their guarantees, the five-stage pipeline, edge buffering and replay, the time-series schema, late and out-of-order data, twin state synchronization, failure modes, and a practical checklist.

Context and Background

Plant data reaches an enterprise through a layered stack that has barely changed in two decades: controllers at the bottom, an OPC UA or vendor server exposing tags, a historian or SCADA layer that stores them, and everything above that consuming from the historian. Digital twins arrive and ask for something the historian was not designed for. They want the same signal at three granularities, with explicit quality, joined to a model of the asset, and with changes pushed rather than polled.

Three source families dominate. OPC UA servers expose a typed address space and, through the Subscription service defined in OPC 10000-4, push value changes to clients on a publishing interval, with per-item queues that can discard oldest or newest values on overflow and flag the loss in the status code. MQTT brokers, increasingly organized as a unified namespace, carry the same data as retained topic state plus a stream of changes; Sparkplug B adds birth and death certificates and sequence numbers on top. Classic historians and PLC-side data loggers hold the long tail, accessible through OPC UA HistoryRead or vendor APIs. Our comparison of OPC UA FX, MQTT Sparkplug B and the unified namespace covers the protocol trade-offs; this article assumes you have chosen and asks what happens next.

The standard references are worth reading directly rather than through vendor summaries. The OPC Foundation publishes the OPC UA specification online, and the MQTT 5.0 specification from OASIS defines the session, QoS and expiry semantics that every buffering decision in this article leans on.

The thesis here is narrow and, I think, underused: a digital twin needs two different stores with two different consistency rules. The time-series store is an append-mostly record of observations, and it should accept late and duplicate data gracefully. The twin state is a small, versioned, current-belief view, and it should reject stale updates. Treating them as one thing is the root of most ingestion bugs.

A Reference Architecture for Digital Twin Data Ingestion

A robust ingestion path has five stages: acquire at the source with source timestamps, buffer durably at the edge, transport with at-least-once delivery, normalize and deduplicate on arrival, then fan out to a history store and a twin state service. Each stage owns one guarantee, and none may silently weaken the guarantee of the stage before it.

digital twin data ingestion reference architecture from OPC UA through MQTT to time-series store and twin state

Figure 1: Reference architecture for digital twin data ingestion. Data flows left to right, with a durable buffer at the edge and a split between history and state at the end.

Figure 1 shows the reference path. An OPC UA server fronts the controller; an edge gateway subscribes, stamps and buffers; a broker in the unified namespace carries data to consumers; and the consumers write to a raw time-series store and to a twin state service. The twin model is built from both: history-derived features such as rollups, and the current state. Applications and simulations read the twin, not the raw store.

Stage one: acquire with source time, not arrival time

The single most important field in the pipeline is the timestamp assigned closest to the physical event. OPC UA gives you two: SourceTimestamp, set by the device or the server’s best knowledge of when the value was sampled, and ServerTimestamp, set when the server processed it. Prefer the source timestamp and keep both. If a PLC does not provide one, the server stamps it, and you have a weaker guarantee that you should record rather than hide.

Do not stamp on arrival at the gateway unless you must. Arrival time conflates sampling time with network and buffering delay, which is exactly the quantity you need to separate when a link has been down for six hours and replays. Keep a third time, the ingest time, in the database; the gap between event time and ingest time is your free, always-on lateness monitor.

Subscribe rather than poll. A subscription pushes changes at a negotiated publishing interval, and each monitored item has its own sampling interval and queue. Setting queue size greater than one lets the server hold several samples between publishes, so a 100 ms signal is not flattened to the publishing rate. When that queue overflows, the specification has the server discard according to a configured policy and set the overflow flag in the DataValue status code, so a consumer that checks status can know it lost samples rather than assume it did not.

Stage two: buffer durably at the edge

The gateway is where store-and-forward IIoT earns its keep. Plant-to-cloud links fail routinely: a firewall change, a modem reboot, a maintenance window. If the gateway holds samples only in memory, every outage becomes a data hole, and the twin quietly diverges from reality for the duration.

The buffer should be an append-only local log on disk with a monotonically increasing sequence number per gateway, a record of acknowledgment per entry, and a bounded retention policy expressed in bytes and hours. SQLite in WAL mode, an embedded log such as a segmented file queue, or a small broker on the gateway all work. What matters is that a sample is on disk before it is acknowledged upward and is removed only after the next stage acknowledges it.

Stage three: transport with at-least-once semantics

MQTT 5.0 gives you the tools. QoS 1 provides at-least-once delivery, meaning duplicates are possible and must be tolerated downstream. QoS 2 provides exactly-once delivery between a client and broker, at the cost of a four-packet handshake and held state; in practice most industrial pipelines choose QoS 1 plus an idempotent consumer, because the end-to-end exactly-once guarantee has to be built at the application layer anyway.

A persistent session is the second tool. With a clean start flag of false and a Session Expiry Interval set, the broker retains the session, including queued QoS 1 messages for subscribers and subscriptions, across disconnects until the interval lapses. The interval is a four-byte count of seconds; zero or absent ends the session at disconnect, and 0xFFFFFFFF means it never expires. Choose it deliberately: a value shorter than your longest planned consumer outage silently turns a restart into data loss.

Stage four: normalize and deduplicate on arrival

The ingest consumer converts heterogeneous payloads into one canonical record: signal identity, event time, ingest time, value, quality, and provenance. Identity deserves care. A signal should have a stable internal identifier that survives re-tagging in the PLC, with the OPC UA NodeId, the MQTT topic and the vendor tag name stored as aliases. When a tag is renamed during a commissioning change, the history should follow the signal, not the string.

Deduplication happens here by a natural key: signal identifier plus event time, optionally plus a gateway sequence. That key is what lets the transport stay at-least-once without polluting history, and we will make it concrete in the schema section below.

Stage five: fan out to history and state

The same normalized record feeds two writers. The history writer appends to the time-series store without any notion of “current.” The state writer updates the twin’s current view only if the record is newer than what it already holds. They share an input and nothing else, and keeping them independent means a slow state service never back-pressures history capture.

The rest of the article takes these stages in turn, with code, and then returns to the failure modes that cut across them.

Walk-through: Subscribing, Buffering, Publishing, and Replaying

The practical path is an OPC UA subscription that writes every change to a durable local log, a publisher that drains the log to MQTT with QoS 1, and a replay mode that marks old samples as historical. The code below is a minimal, runnable skeleton in Python using asyncua, SQLite and paho-mqtt; production gateways add batching, TLS and certificate handling.

Subscribing with a queue and capturing source time

The asyncua library creates a subscription with a requested publishing period in milliseconds and lets you pass a handler. The handler receives the node, the value, and the data change notification, from which the monitored item’s DataValue gives the source timestamp and status.

import asyncio, sqlite3, json
from asyncua import Client

DB = sqlite3.connect("buffer.db", isolation_level=None)
DB.execute("PRAGMA journal_mode=WAL")
DB.execute("""CREATE TABLE IF NOT EXISTS outbox(
  seq INTEGER PRIMARY KEY AUTOINCREMENT,
  signal TEXT NOT NULL,
  src_ts TEXT NOT NULL,
  value TEXT NOT NULL,
  status INTEGER NOT NULL,
  acked INTEGER NOT NULL DEFAULT 0)""")

class Handler:
    def datachange_notification(self, node, val, data):
        dv = data.monitored_item.Value
        ts = (dv.SourceTimestamp or dv.ServerTimestamp).isoformat()
        DB.execute(
            "INSERT INTO outbox(signal, src_ts, value, status) VALUES (?,?,?,?)",
            (node.nodeid.to_string(), ts, json.dumps(val),
             dv.StatusCode.value))

async def main():
    url = "opc.tcp://plc-gw.plant.local:4840"
    async with Client(url=url) as client:
        nodes = [client.get_node("ns=2;s=Pump01.Flow"),
                 client.get_node("ns=2;s=Pump01.Pressure")]
        sub = await client.create_subscription(500, Handler())
        await sub.subscribe_data_change(nodes, queuesize=10)
        await asyncio.sleep(10**9)

asyncio.run(main())

Two details matter. The insert happens before anything is published, so a crash after the subscription callback loses at most the in-flight notification the OPC UA stack had not yet delivered. And status is stored as a number, so a Bad or Uncertain sample is preserved rather than dropped; the twin should decide what quality means, not the gateway.

A caution on completeness: a subscription gives you changes since you subscribed. If the gateway process is down, the server does not retain notifications for you indefinitely, since subscription lifetime is bounded by negotiated counts. For gaps longer than that, use OPC UA HistoryRead against the server’s own history, if it has one, to backfill. Plan this path before you need it.

Publishing with QoS 1 and a persistent session

The publisher reads unacknowledged rows in sequence order and marks them after the broker’s PUBACK. With paho-mqtt version 2, the callback API version is explicit and MQTT 5 properties are objects.

import json, sqlite3, time
import paho.mqtt.client as mqtt
from paho.mqtt.packettypes import PacketTypes
from paho.mqtt.properties import Properties

DB = sqlite3.connect("buffer.db", isolation_level=None)
inflight = {}   # mid -> seq

def on_publish(client, userdata, mid, reason_code, properties):
    seq = inflight.pop(mid, None)
    if seq is not None:
        DB.execute("UPDATE outbox SET acked=1 WHERE seq=?", (seq,))

c = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2,
                client_id="gw-line3-01", protocol=mqtt.MQTTv5)
c.on_publish = on_publish
c.tls_set()                       # use real certificates in production

props = Properties(PacketTypes.CONNECT)
props.SessionExpiryInterval = 86400   # seconds; keep session one day
c.connect("broker.plant.example", 8883, keepalive=30,
          clean_start=False, properties=props)
c.loop_start()

while True:
    rows = DB.execute(
        "SELECT seq, signal, src_ts, value, status FROM outbox "
        "WHERE acked=0 ORDER BY seq LIMIT 200").fetchall()
    for seq, signal, ts, value, status in rows:
        body = json.dumps({"s": signal, "t": ts, "v": json.loads(value),
                           "q": status, "gw": "gw-line3-01", "n": seq})
        info = c.publish(f"plant1/line3/pump01/{signal}", body, qos=1)
        inflight[info.mid] = seq
    time.sleep(0.2 if rows else 1.0)

Note what is absent: the publisher never deletes on send. It marks on acknowledgment, and a separate janitor purges acknowledged rows older than a retention horizon. Note also what is a trap: the loop above re-queries unacknowledged rows each pass, so un-acked rows in flight are republished. That is deliberate at-least-once behavior, and it is exactly why the consumer must deduplicate. A production version tracks in-flight sequence numbers to avoid needless repeats and respects the broker’s Receive Maximum, the limit on concurrent QoS 1 and 2 publications, which defaults to 65,535 if the broker does not set it.

store and forward IIoT sequence: gateway buffers during a WAN outage and replays flagged historical samples on reconnect

Figure 2: Store-and-forward sequence. The gateway keeps appending while the link is down, then replays in order, and the consumer deduplicates on signal and event time.

Figure 2 traces the loop through an outage. While the link is up, each change is appended, published and marked acknowledged. When the link drops, appends continue and nothing is marked. On return the gateway replays in sequence order; because delivery is at-least-once, the consumer may see some samples twice and discards the repeats by key.

Replay semantics: live versus historical

A replayed sample is not a live sample, and consumers need to tell the difference. A live sample should update the twin state; a replayed sample from six hours ago should update history only, unless it is newer than the current state. Sparkplug B has a field for this: a metric in a payload can carry an is_historical flag, which the specification defines so that data buffered while a node was offline is not mistaken for a current value. If you are not using Sparkplug, add an explicit field in your envelope.

Also decide the replay order. Replaying oldest first preserves causal order and is simple to reason about, but it delays fresh data behind a large backlog after a long outage. A common compromise is to interleave: send the newest sample per signal immediately so state heals fast, then drain the backlog oldest-first at a capped rate. This works only because state and history are separate writers, which is the payoff of the earlier split.

Sizing the buffer

Size the buffer from first principles rather than hope. A signal sampled at 1 Hz with a 100-byte serialized envelope produces about 8.6 MB per day; two thousand such signals produce roughly 17 GB per day before compression. These are arithmetic illustrations, not measurements of any product. If your design target is a 24-hour outage with headroom, that gateway needs tens of gigabytes of retained log, a number that surprises teams who sized the edge box for compute alone.

Decide in advance what is dropped when the buffer fills, and tell the twin. Dropping oldest preserves the most recent state and is the usual choice; dropping newest preserves a contiguous history but starves state. Whichever you choose, emit a gap marker with the time range lost, so the downstream store records a known hole, not a silent one. We cover the broker side of this pattern in the MQTT Sparkplug B reference architecture.

The Time-Series Schema: Shape, Keys and Idempotency

Design the history store around a narrow long table keyed by signal and event time, with a unique constraint that makes inserts idempotent. Keep value columns typed rather than JSON, keep quality explicit, and keep both event time and ingest time. Everything else, rollups included, is derived.

The table below is an example in PostgreSQL with the TimescaleDB extension. The same logical design maps onto ClickHouse or InfluxDB; our InfluxDB, TimescaleDB and ClickHouse comparison explains when to choose which.

CREATE TABLE signal (
  signal_id   BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
  asset_id    TEXT NOT NULL,
  name        TEXT NOT NULL,
  unit        TEXT,
  value_type  TEXT NOT NULL CHECK (value_type IN ('float','int','bool','text')),
  opcua_node  TEXT,
  mqtt_topic  TEXT,
  UNIQUE (asset_id, name)
);

CREATE TABLE sample (
  signal_id   BIGINT      NOT NULL REFERENCES signal(signal_id),
  event_ts    TIMESTAMPTZ NOT NULL,
  ingest_ts   TIMESTAMPTZ NOT NULL DEFAULT now(),
  v_double    DOUBLE PRECISION,
  v_text      TEXT,
  quality     SMALLINT    NOT NULL DEFAULT 0,
  gw_id       TEXT,
  gw_seq      BIGINT,
  PRIMARY KEY (signal_id, event_ts)
);

SELECT create_hypertable('sample', 'event_ts',
                         chunk_time_interval => INTERVAL '1 day');

ALTER TABLE sample SET (
  timescaledb.compress,
  timescaledb.compress_segmentby = 'signal_id',
  timescaledb.compress_orderby   = 'event_ts DESC');
SELECT add_compression_policy('sample', INTERVAL '7 days');

The composite primary key does the heavy lifting. TimescaleDB requires that unique indexes on a hypertable include the partitioning column, and event_ts is already part of this one, so it is permitted. The ingest statement then becomes a one-liner that is safe under replay:

INSERT INTO sample (signal_id, event_ts, v_double, quality, gw_id, gw_seq)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (signal_id, event_ts) DO NOTHING;

Why event time plus signal is the right natural key

A gateway sequence number alone is not enough, because sequences reset when a gateway is replaced, and one physical event can be seen by two redundant gateways. Event time per signal is a stable identity for an observation, with one honest weakness: two distinct samples at the same instant collide. At millisecond or finer resolution that is rare for a single signal, and you should treat a collision with different values as a data-quality event, not silently drop it. An alternative is DO UPDATE when quality improves, which is useful when a Bad sample is later superseded by a Good one.

Typed columns and what to do with strings

Store numeric values as double precision and booleans as 0 or 1 in the same column unless your analytics need native booleans. Strings and enumerations go in a separate column or table, because compressing numeric columns is where columnar and delta encodings pay off. Avoid a catch-all JSON value: it is easy to write and expensive to query, and it hides schema drift until a dashboard breaks.

Rollups and the dirty-bucket pattern

Twins rarely want raw 10 Hz data. They want one-minute and one-hour aggregates, and they want them to be correct after late data arrives. Continuous aggregates in TimescaleDB, materialized views in ClickHouse, or a scheduled job can all do this, but the design question is the same: when a late sample lands in a bucket that was already computed, how does the system know to recompute it? The next section answers that.

Late and Out-of-Order Data Without Corrupting the Twin

Late data is not an edge case in industrial systems; it is the normal result of buffering, which you just built on purpose. The ingest path needs an explicit policy for three situations: samples that arrive after their rollup was computed, samples with timestamps that are implausible, and samples that arrive in a different order than they happened.

time series ingestion architecture decision flow for late and out-of-order data with dirty rollup buckets

Figure 3: Ingest decision flow. Invalid time goes to a dead letter queue, valid samples go to a hot or cold chunk, and late arrivals mark rollup buckets dirty for recomputation.

Figure 3 shows the flow. The first gate validates event time, rejecting timestamps far in the future or before the plausible commissioning date of the asset. A PLC with a dead clock battery reporting 1970 should land in a dead letter table with its raw payload, where an engineer can see it, not in your history as a 55-year-old measurement. The second gate asks whether the sample falls inside the hot window, which for most systems means the chunks still uncompressed and cheap to write. Inside the window, an insert is routine. Outside it, the write is still allowed, but the affected rollup bucket is marked dirty.

The dirty-bucket pattern

A dirty-bucket table records (signal_id, bucket_start, resolution) for every late write. A background job recomputes each dirty bucket from raw samples, upserts the aggregate, and clears the mark. The cost is proportional to late data volume, not history size, which makes it safe to leave running.

CREATE TABLE rollup_dirty (
  signal_id    BIGINT      NOT NULL,
  bucket_start TIMESTAMPTZ NOT NULL,
  resolution   INTERVAL    NOT NULL,
  marked_at    TIMESTAMPTZ NOT NULL DEFAULT now(),
  PRIMARY KEY (signal_id, bucket_start, resolution)
);

-- recompute one dirty hourly bucket
INSERT INTO rollup_1h (signal_id, bucket_start, v_avg, v_min, v_max, n)
SELECT s.signal_id, d.bucket_start,
       avg(v_double), min(v_double), max(v_double), count(*)
FROM rollup_dirty d
JOIN sample s
  ON s.signal_id = d.signal_id
 AND s.event_ts >= d.bucket_start
 AND s.event_ts <  d.bucket_start + d.resolution
WHERE d.resolution = INTERVAL '1 hour'
  AND s.quality = 0
GROUP BY s.signal_id, d.bucket_start
ON CONFLICT (signal_id, bucket_start) DO UPDATE
SET v_avg = EXCLUDED.v_avg, v_min = EXCLUDED.v_min,
    v_max = EXCLUDED.v_max, n = EXCLUDED.n;

Note the n column. Storing the sample count with every aggregate lets consumers see how much evidence supports an hourly average, and it makes a bucket built from 12 of an expected 3,600 samples visibly suspicious. Completeness is a first-class property of a rollup, not an afterthought.

Telling consumers a correction happened

When a bucket is recomputed, anything that already consumed the old value is now stale: a feature in a model, a KPI snapshot, a report. The last box in Figure 3 notifies the twin. A lightweight approach is to publish a correction event on a dedicated topic carrying the signal, the interval and a new version number, so downstream computations can invalidate and recompute. Whether to recompute everything is a business decision; the architectural point is that corrections are visible, not silent.

Out-of-order within the hot window

Within a short window, samples from different gateways, or from replay interleaved with live traffic, will interleave. For history this is harmless, because the key makes the insert order-independent. For derived streaming computations such as a moving average, it is not, and you need either a watermark that delays emission by the maximum expected lateness, or a recompute-on-correct design. A watermark of a few seconds suits a healthy network; it is the wrong tool for a six-hour replay, which is a correction case.

Clock discipline

Everything above assumes event timestamps are comparable across sources. They are only as good as the clocks. OPC UA servers and gateways should discipline time against NTP or, for tighter needs, PTP, and the gateway should record its own clock offset in the envelope when it can measure one. Drift of even a second between two redundant gateways turns the dedupe key into a near-miss that stores both copies. A cheap guard is a nightly query for signals with unexpectedly high sample counts per bucket, which usually means duplicates with skewed clocks.

Twin State Synchronization: History Is Not State

A twin’s current state is a derived, versioned view, updated only by samples newer than what it already holds. History accepts everything, including old data; state accepts only progress. Separating them, with a merge patch applied under a newer-wins rule and a version stamp on every change, keeps replays from rolling the twin backwards.

twin state synchronization flow mapping signals to twin properties with newer-than check and versioned change events

Figure 4: Twin state synchronization. A normalized sample is mapped to a twin property, applied only if newer than stored state, validated, versioned and published.

Figure 4 follows one sample into the twin. First the signal is mapped to a property of a twin node. That mapping is itself part of the model: a pump twin has a flowRate property bound to a particular signal identity, and the binding is data, not code. How you express the model, whether as an Asset Administration Shell submodel, a DTDL interface or an OPC UA information model, is a separate decision, explored in our AAS, DTDL and OPC UA information model comparison.

Newer-wins with a tie-breaker

The rule is simple: apply an update to a property only if its event time is later than the property’s last-applied event time. The edge cases are where designs break. Equal timestamps need a deterministic tie-break, typically the higher gateway sequence or the better quality. A sample with Bad quality should not overwrite a Good value as “current,” but it should be recorded, and the property should be marked stale after a staleness horizon rather than keep presenting an old Good value as if it were live.

def apply_sample(state, prop, sample):
    cur = state.get(prop)
    newer = (cur is None or
             (sample.event_ts, sample.gw_seq) > (cur["event_ts"], cur["gw_seq"]))
    if not newer:
        return None                       # keep in history only
    if sample.quality != 0 and cur is not None and cur["quality"] == 0:
        cur["stale_after"] = sample.event_ts   # remember a degraded reading
        return None
    new = {"value": sample.value, "event_ts": sample.event_ts,
           "gw_seq": sample.gw_seq, "quality": sample.quality,
           "version": (cur["version"] + 1) if cur else 1}
    state[prop] = new
    return new                            # emit as a change event

This is a sketch, not a framework; a real service persists state transactionally and makes the compare-and-set atomic, for example with an optimistic version check or a single-row UPDATE ... WHERE event_ts < $new statement.

Staleness is information

A twin that never admits ignorance is dangerous. Give every property a freshness field derived from the expected update period of its signal: a 1 Hz pressure that has not updated for 30 seconds is stale, while a daily lab result is not. Expose staleness through the same API as the value, so a simulation can decide to run on the last known value, an estimate, or refuse. Unified namespace deployments help here because retained MQTT state and birth and death certificates already convey liveness; see the unified namespace architecture primer for the topic design.

Initial state and reconnect

When a state service starts, it has no memory. It can rebuild from the last known value per signal in the history store, from retained messages on the broker, or, with Sparkplug, by requesting a rebirth so edge nodes republish their full set of metrics. Retained messages are convenient but have one trap: they hold the last value indefinitely, so a retained topic from a decommissioned device looks alive forever unless something clears it. Prefer the history store as the source of truth for rebuilds and use retained state as a fast cache.

Versioned change events

Every applied change gets a version, and the state service publishes it. Version numbers per property let subscribers detect gaps, such as receiving version 41 after 39, and request a resync. They also provide a clean audit answer to the question a regulated plant will ask: what did the twin believe at 14:03, and what evidence did it have? Because history keeps every observation and state keeps versioned beliefs, the answer is reconstructable.

Trade-offs, Gotchas, and What Goes Wrong

The failures below are the ones that recur across plants. None is exotic; each follows from a guarantee that one stage assumed and another did not provide.

Subscription gaps are invisible unless you look. OPC UA subscriptions can lose samples to queue overflow, and the server flags this in the status code. If your gateway discards status, you have turned a detectable loss into an undetectable one. Preserve status, and alarm on overflow flags.

Counting on exactly-once. QoS 2 does not give you exactly-once processing. A consumer that crashes after receiving but before committing will see the message again after reconnect with a persistent session. The cure is the idempotent key, not a higher QoS.

Session expiry shorter than the outage. A session expiry interval of an hour with an eight-hour consumer outage drops queued messages for that consumer. Pair the interval with an alert on consumer lag, and remember that brokers also bound queue depth by their own configuration, which is vendor-specific and should be verified for your broker.

Edge disk fills silently. A buffer without a bound and a gap marker turns a long outage into either a crashed gateway or an unannounced hole. Test it: pull the network cable on a staging gateway for twice the planned outage and watch what happens.

Replay storms. After a long outage, many gateways reconnect together and replay at full speed, saturating the broker and the database. Rate-limit replay per gateway and give live data priority, as described earlier.

Cardinality and write amplification. Thousands of signals at high frequency in a naive row-per-sample design create index and storage costs that dominate the project. Report-by-exception and deadband filtering at the gateway reduce volume dramatically for slow-moving signals, but they change the meaning of a missing sample: the absence of data now means “no change,” which the state layer must interpret correctly, usually through a heartbeat sample at a minimum interval.

Schema drift at the source. A firmware update changes an engineering unit from bar to kPa. Without unit metadata in the signal table and a change-detection step, the twin silently mixes units. Store units, version the signal definition, and treat a change as a new signal version.

Security is an ingestion concern. The edge gateway holds credentials to a control network and to the cloud. Use OPC UA SecurityPolicy with signed and encrypted sessions, certificate trust lists, and MQTT over TLS with per-gateway identities and topic-level authorization. A buffer on disk contains plant data in plaintext unless you encrypt it.

Over-modeling too early. Teams spend months on an elegant ontology before a single reliable signal reaches the store. Ingest raw, well-keyed data first; the history table does not care what the twin model eventually becomes.

Practical Recommendations

Build the pipeline in the order of its guarantees: durable capture first, idempotent storage second, state sync third, modeling last. Measure each stage with a small set of numbers, such as ingest lag, sample completeness per signal and replay backlog, and treat any unmeasured stage as unverified.

A sensible default stack for a first production line is an OPC UA subscription client on an edge gateway with a disk-backed outbox, MQTT 5 with QoS 1 and persistent sessions into a broker hosting the unified namespace, an idempotent consumer writing to a hypertable or equivalent columnar store, a separate state service with versioned properties, and a dead letter path that someone actually reads. Alternatives are legitimate: Kafka in place of MQTT for high-volume plant-to-cloud links, or a managed cloud ingestion service in place of a self-run consumer. The guarantees must still be assigned to a stage.

Checklist before you call the pipeline production-ready:

  • Source timestamps captured and stored beside ingest time; both queryable.
  • Status or quality preserved end to end, with overflow flags alarmed.
  • Edge buffer on disk, bounded, with a documented drop policy and gap markers.
  • QoS 1 with persistent sessions; expiry interval longer than the longest tolerated outage.
  • Unique key on signal and event time; replay tested by publishing every message twice.
  • Replay rate-limited, flagged as historical, and kept out of state unless newer.
  • Dirty-bucket recompute for rollups, with sample counts stored.
  • State updates newer-wins, versioned, with explicit staleness.
  • Clock offsets monitored; dead letter queue reviewed weekly.
  • A tested rebuild procedure for the state service from history.

If you do only one thing from this list, run the duplicate-publish test. It exposes the largest share of ingestion bugs for the least effort, and passing it is the clearest evidence that your pipeline is truly at-least-once safe.

Frequently Asked Questions

What is digital twin data ingestion?

Digital twin data ingestion is the pipeline that acquires signals from machines and systems, buffers and transports them, normalizes them into a canonical record and delivers them to a history store and a twin state service. It is distinct from modeling: ingestion decides what the twin can know and how fresh and trustworthy that knowledge is. Its core guarantees are source timestamps, durable buffering, at-least-once delivery, idempotent storage and versioned state.

Should I use OPC UA subscriptions or MQTT for the twin?

They are usually complementary, not competing. OPC UA subscriptions are the standard way to acquire typed, timestamped values from a controller or server, with per-item queues and status codes. MQTT is a good transport and distribution layer for fan-out to many consumers, especially in a unified namespace. A common pattern is OPC UA at the edge for acquisition and MQTT from the gateway upward, with source timestamps and quality carried in the payload.

How do I handle network outages without losing data?

Persist every sample to a disk-backed outbox on the edge gateway before acknowledging it upward, publish with MQTT QoS 1, and remove entries only after the broker acknowledges them. On reconnect, replay in order at a capped rate and flag replayed samples as historical. Size the buffer from signal count, rate and the longest outage you tolerate, and define a drop policy with gap markers for when it fills.

How do I deduplicate at-least-once data?

Define a natural key, typically signal identifier plus event timestamp, and enforce it with a unique constraint so repeated inserts are no-ops. In PostgreSQL or TimescaleDB, INSERT ... ON CONFLICT DO NOTHING does this. Treat a conflict with a different value as a data-quality event rather than silently discarding it, and verify the behavior by deliberately publishing every message twice in a test environment.

What is the difference between history and twin state?

History is an append-only record of every observation, accepting late and duplicate data. State is the twin’s current belief per property, updated only when a sample is newer than what is stored, and versioned so subscribers can detect gaps. Keeping them separate prevents a replay of old data from rolling the twin backwards and lets you reconstruct what the twin believed at any moment.

How should I deal with late data in rollups?

Mark affected rollup buckets as dirty whenever a late sample is written, and recompute them from raw data in the background. Store the sample count with each aggregate so completeness is visible, and publish a correction event so downstream computations can invalidate dependent results. Streaming computations can use a watermark for small lateness, but long replays should be treated as corrections rather than ordinary lateness.

Further Reading

By Riju — about

Comments

No comments yet. Why don’t you start the discussion?

Leave a Reply

Your email address will not be published. Required fields are marked *