Real-Time Aptos Transaction Streaming
Tail the head of the Aptos chain with aptos.indexer.v1.RawData/GetTransactions. Covers batch handling, server-side filters, checkpointing from processed_range, and reconnect strategy.
Early access — not yet in production
The Aptos Transaction Stream is not generally available, and we have not measured delivery lag at the head of chain. Email support@dwellir.com to join the early-access list.
A real-time tail subscribes near the head of the chain and receives finalized transactions as blocks are produced. It replaces polling loops in indexers, notification services, dashboards, and trading systems.
Read the Transaction Stream Overview first — it covers the service definition, request fields, and metering that this page builds on.
Start a Tail
Omit starting_version to begin near the head of chain. Omit transactions_count
so the stream runs indefinitely.
import { credentials, Metadata } from '@grpc/grpc-js';
const client = new RawDataClient(
process.env.DWELLIR_APTOS_STREAM_ENDPOINT, // assigned at onboarding
credentials.createSsl()
);
const metadata = new Metadata();
metadata.add('x-api-key', process.env.DWELLIR_API_KEY);
// No starting_version -> begin at the head of chain
// No transactions_count -> unbounded stream
const request = {};
const stream = client.GetTransactions(request, metadata);
stream.on('data', (response) => {
// Each message is a BATCH, not a single transaction.
for (const tx of response.transactions) {
handleTransaction(tx);
}
// Advance progress from processed_range, not from the last transaction.
if (response.processed_range) {
lastProcessedVersion = BigInt(response.processed_range.last_version);
}
});
stream.on('error', (error) => {
console.error('Stream error:', error.code, error.details);
scheduleReconnect();
});
stream.on('end', () => {
scheduleReconnect();
});One message is a batch of transactions
TransactionsResponse carries transactions[], chain_id, and
processed_range. A handler written as stream.on('data', tx => tx.version)
reads the batch as though it were a transaction and will silently process nothing
useful. Always iterate response.transactions.
Read a Transaction
Top-level fields come from aptos.transaction.v1.Transaction. The hash lives on
info and is a bytes field:
function handleTransaction(tx) {
const hash = Buffer.from(tx.info.hash).toString('hex');
console.log({
version: tx.version, // string (jstype = JS_STRING)
hash: `0x${hash}`,
blockHeight: tx.block_height,
epoch: tx.epoch,
type: tx.type, // TRANSACTION_TYPE_USER, _BLOCK_METADATA, ...
success: tx.info.success,
gasUsed: tx.info.gas_used,
timestamp: tx.timestamp, // { seconds, nanos }
});
if (tx.type === 'TRANSACTION_TYPE_USER') {
processUserTransaction(tx.user);
}
}version, block_height, epoch, and gas_used are annotated
jstype = JS_STRING upstream, so most JavaScript stubs hand them back as strings.
Convert with BigInt() before doing arithmetic.
Filter on the Server
A head-of-chain tail on Aptos carries every transaction the network produces. Filtering server-side cuts both the data you parse and the messages you are billed for, because the filter runs before delivery and before metering.
// Successful 0x1::coin::transfer calls only.
const request = {
transaction_filter: {
logical_and: {
filters: [
{ api_filter: { transaction_root_filter: { success: true } } },
{
api_filter: {
user_transaction_filter: {
payload_filter: {
entry_function_filter: {
address: '0x1',
module_name: 'coin',
function: 'transfer',
},
},
},
},
},
],
},
},
};BooleanTransactionFilter composes three leaf filters — transaction_root_filter
(success, transaction type), user_transaction_filter (sender, entry function),
and event_filter (Move struct type, data substring) — with logical_and,
logical_or, and logical_not.
A filtered stream still sends progress messages
When a filter removes everything in a version range, the server still emits a
message carrying processed_range so you can advance your checkpoint. Those
messages are metered. A narrow filter reduces bytes and parsing work; it does not
reduce message count to zero.
Track Progress
processed_range reports the version span the server covered, including versions
your filter dropped. Checkpointing from it keeps progress moving even during long
runs of non-matching transactions, and it is what you supply as starting_version
after a restart.
class TailCheckpoint {
private lastVersion: bigint | null = null;
private dirty = false;
observe(response) {
if (!response.processed_range) return;
this.lastVersion = BigInt(response.processed_range.last_version);
this.dirty = true;
}
// Flush on an interval rather than per message.
async flush() {
if (!this.dirty || this.lastVersion === null) return;
await db.upsert('stream_checkpoints', {
name: 'aptos_tail',
last_version: this.lastVersion.toString(),
updated_at: new Date(),
});
this.dirty = false;
}
async resumeVersion(): Promise<bigint | null> {
const row = await db.findOne('stream_checkpoints', { name: 'aptos_tail' });
return row ? BigInt(row.last_version) + 1n : null;
}
}Only checkpoint versions you have fully processed downstream. Recording progress before the write lands turns a crash into a silent gap.
Reconnect Without Gaps
Reconnect with starting_version set to the version after your last checkpoint.
A reconnect can redeliver versions you already handled, so downstream writes must
be idempotent.
class TailManager {
private attempts = 0;
async connect() {
const resumeFrom = await this.checkpoint.resumeVersion();
// First run has no checkpoint: start at the head.
const request = resumeFrom === null
? {}
: { starting_version: resumeFrom.toString() };
const stream = client.GetTransactions(request, metadata);
stream.on('data', (response) => {
this.attempts = 0;
this.handleBatch(response);
});
stream.on('error', () => this.scheduleReconnect());
stream.on('end', () => this.scheduleReconnect());
}
private async scheduleReconnect() {
this.attempts += 1;
const backoff = Math.min(1000 * 2 ** this.attempts, 30_000);
const jitter = Math.random() * 0.3 * backoff;
await sleep(backoff + jitter);
await this.connect();
}
}Aptos reaches BFT finality before a transaction enters the stream, so there is no reorg handling to write. A version you have processed will not be revised.
Handle Backpressure
@grpc/grpc-js emits data faster than most consumers can write. Serialize
processing and pause the stream when your queue grows, rather than letting an
unbounded array absorb the difference.
const MAX_QUEUE = 500;
const queue = [];
let draining = false;
stream.on('data', (response) => {
queue.push(response);
if (queue.length > MAX_QUEUE) {
stream.pause();
}
if (!draining) void drain();
});
async function drain() {
draining = true;
while (queue.length > 0) {
const response = queue.shift();
for (const tx of response.transactions) {
await handleTransaction(tx);
}
checkpoint.observe(response);
if (queue.length < MAX_QUEUE / 2) {
stream.resume();
}
}
draining = false;
}Expected Message Rate
For an unfiltered tail, the server sends roughly one message per block. It does
not wait to fill batch_size, so head-of-chain batches are much smaller than the
1000 ceiling and your message rate tracks block production rather than transaction
volume.
Because the billed unit is the message, this is the number to compare against your existing Sui checkpoint usage when estimating cost. See Metering.
Monitor the Tail
Track these four signals:
| Signal | Why it matters |
|---|---|
| Version lag | Head version minus your last checkpoint. The clearest sign a consumer is falling behind. |
| Messages per second | Compare against block rate to detect a stalled or throttled stream. |
| Time since last message | A silent stream and a healthy quiet period look identical without a threshold. |
| Reconnect rate | Rising reconnects usually mean backpressure or rate limiting, not network faults. |
A tail that stops receiving messages without an error is usually rate limiting. On Dwellir, plan rate limits apply per account, not per key, so a second workload on the same account can throttle this stream — see Rate Limits.
Related Guides
- Transaction Stream Overview — service definition, request fields, metering
- Historical Replay — starting from a stored version
- Custom Processors — worker architecture on batch semantics
- Authentication — API key handling
Transaction Stream Overview
Stream finalized Aptos transactions over gRPC with aptos.indexer.v1.RawData/GetTransactions. Covers the request and response shapes, server-side filtering, batch semantics, metering, and rate-limit scope.
Historical Replay
Replay Aptos transactions from a stored ledger version with GetTransactions. Covers starting_version, bounded ranges, checkpoint resumption, and backfill throughput.