Building Aptos Stream Processors
Architecture patterns for Aptos Transaction Stream processors - batch handling, event decoding, idempotent writes, checkpoint ordering, and failure isolation.
Early access — not yet in production
The Aptos Transaction Stream is not generally available. Email support@dwellir.com to join the early-access list.
A processor subscribes to the Transaction Stream, decodes the transactions it cares about, and writes them into your own schema. This page covers the architecture; see the Transaction Stream Overview for the service definition and Real-Time Streaming for the connection and reconnect mechanics.
The Processing Loop
The unit of work is a TransactionsResponse — a batch — not a transaction. Build
the loop around the batch, because the batch is also the checkpoint boundary and
the billing unit.
class AptosStreamProcessor {
async handleBatch(response) {
// Nothing matched the filter in this range, but progress still advanced.
if (response.transactions.length === 0) {
return this.checkpoint.observe(response);
}
const records = [];
for (const tx of response.transactions) {
const decoded = this.decode(tx);
if (decoded) records.push(decoded);
}
// One transactional write per batch, then checkpoint.
await this.db.transaction(async (trx) => {
await trx.bulkUpsert('transactions', records);
await trx.upsert('stream_checkpoints', {
name: 'aptos',
last_version: response.processed_range.last_version,
});
});
}
}Writing the records and the checkpoint in the same database transaction is what makes restarts safe. Split them and a crash between the two either loses versions or reprocesses them.
Empty batches are normal
A filtered stream sends messages with no transactions when nothing in a version
range matched. They still carry processed_range and they are still metered.
Treat an empty batch as a progress update, not as an error or an idle stream.
Filter Server-Side First
Anything you can express as a BooleanTransactionFilter should go in the request
rather than in shouldProcess(). The filter runs on the node, before delivery and
before metering, so server-side filtering reduces your bill while client-side
filtering does not.
// Prefer this...
const request = {
transaction_filter: {
api_filter: {
user_transaction_filter: {
payload_filter: {
entry_function_filter: { address: '0x1', module_name: 'coin' },
},
},
},
},
};
// ...over receiving everything and discarding it in the handler.Client-side filtering still has a place for conditions the filter cannot express — decoded field values, cross-transaction state, or thresholds computed at runtime.
Decoding Events
Events hang off the transaction payload. Event.type_str carries the fully
qualified Move struct tag, and data is a JSON string:
function decodeEvents(tx) {
const events = tx.user?.events ?? tx.block_metadata?.events ?? [];
return events.flatMap((event) => {
const data = JSON.parse(event.data);
switch (event.type_str) {
case '0x1::coin::WithdrawEvent':
return [{
kind: 'coin_withdraw',
amount: BigInt(data.amount),
account: event.key?.account_address,
sequence: BigInt(event.sequence_number),
version: BigInt(tx.version),
}];
case '0x1::coin::DepositEvent':
return [{
kind: 'coin_deposit',
amount: BigInt(data.amount),
account: event.key?.account_address,
sequence: BigInt(event.sequence_number),
version: BigInt(tx.version),
}];
default:
return [];
}
});
}Numeric fields annotated jstype = JS_STRING arrive as strings in most JavaScript
stubs. Parse them with BigInt rather than Number — Aptos amounts in Octas
exceed Number.MAX_SAFE_INTEGER.
Idempotent Writes
A reconnect can redeliver versions you already processed, so every write must be
safe to repeat. Key on version — it is unique and monotonic across the chain:
INSERT INTO coin_transfers (version, event_index, sender, amount)
VALUES ($1, $2, $3, $4)
ON CONFLICT (version, event_index) DO NOTHING;Derived aggregates need the same care. Recompute them from the stored rows, or guard the increment with the version, rather than blindly adding on each delivery.
Isolating Failures
A malformed record should not stall the stream. Route repeated failures to a dead letter table and keep the checkpoint moving:
async function processWithIsolation(tx) {
try {
await processTransaction(tx);
} catch (error) {
await deadLetter.insert({
version: tx.version,
error: error.message,
received_at: new Date(),
});
// Do not rethrow - one bad transaction must not block the batch.
}
}Distinguish the two failure classes. A decode error is specific to one transaction and belongs in the dead letter table. A database outage affects every transaction, and retrying the batch is correct — swallowing it silently loses data.
No reorg handling required
Aptos reaches BFT finality before a transaction enters the stream. Versions are never revised, so processors need no reorg or rollback logic. Do not port confirmation-depth patterns from Ethereum indexers.
Scaling the Write Path
The stream is ordered and single-connection, so parallelism belongs downstream of the reader, not in extra subscriptions.
- Keep one reader. Opening several tails of the same range multiplies your metered messages without adding throughput.
- Parallelize by partition. Shard decoded records across workers by account or event type, then checkpoint only once every shard has committed its slice.
- Batch database writes to whole batches rather than per transaction.
- Flush checkpoints on an interval, not on every message. Checkpoint writes are small and frequent, and they dominate write load if left unbatched.
Monitoring
| Metric | What it tells you |
|---|---|
| Version lag | Head version minus checkpoint. The headline health signal. |
| Batch processing time | Rising time means you will fall behind before lag shows it. |
| Dead letter rate | A spike usually means a schema change upstream. |
| Messages per second | Compare against your plan's limit and your Sui usage. |
Related Guides
- Transaction Stream Overview — service definition and metering
- Real-Time Streaming — connection and reconnect mechanics
- Historical Replay — backfilling a new processor
- Authentication — API key handling
Historical Replay
Replay Aptos transactions from a stored ledger version with GetTransactions. Covers starting_version, bounded ranges, checkpoint resumption, and backfill throughput.
Authentication
How to authenticate against Dwellir's Aptos endpoints - API key in the URL path for REST and GraphQL, x-api-key gRPC metadata for the Transaction Stream.