Stream Ingresses
A stream ingress is an ingress that consumes a retained, ordered event stream through a connection and invokes a reactor component for each event. Where a queue ingress consumes work messages once, a stream ingress replays events: every consumer reads the stream from its own durable position, and the platform persists that position as checkpoints so processing resumes where it left off after any restart.
INFO
The Azure Event Hubs connection type is the initial stream source. RabbitMQ and Azure Service Bus are queue semantics and do not support stream ingresses.
Defining a stream ingress
A stream ingress can be created through the HTTP API, applied as a CLI manifest file, or declared in an app's manifest — whichever fits how you work:
PUT /resources/ingresses/events
Content-Type: application/json{
"componentId": "ingresses/events.hrc",
"type": "stream",
"properties": {
"connectionId": "apps/my-app/partner-events",
"address": "orders-events",
"priority": 100,
"subscription": {
"id": "order-processor",
"startPosition": "earliest",
"checkpoint": { "interval": 10 }
},
"poisonMessage": { "afterAttempts": 5, "action": "quarantine" },
"body": { "mode": "structured" },
"applicationProperties": { "eventType": "eventType" }
},
"acl": ["actors/ticket:create:*"]
}Properties
| Property | Required | Description |
|---|---|---|
connectionId | Yes | The connection to consume from. |
address | Yes | The retained stream address, interpreted by the adapter (e.g. an Event Hubs name). |
priority | No | Collision priority for the same connection+address+subscription; higher wins. |
subscription.id | Yes | The durable reader identity. Adapters map it to their own independent-reader mechanism — an Event Hubs consumer group. Different subscriptions are independent consumers and never collide. |
subscription.startPosition | No | Where a reader with no stored checkpoint starts: earliest (default) or latest. This is not a rewind setting — move existing consumers with the checkpoint management action. |
subscription.checkpoint.interval | No | Minimum seconds between persisted checkpoint advances (default 10). Success is still recorded at-least-once per partition; the interval bounds write frequency. |
poisonMessage.afterAttempts | No | Delivery attempts before the poison policy applies (default 5), so one event cannot permanently stall a partition. |
poisonMessage.action | No | quarantine (default and only action): stop retrying the event, advance the checkpoint past it, and surface it as an error record and a signal. It is never silently dropped. |
body | No | Body binding; same modes as a queue ingress. |
applicationProperties | No | Maps named message application properties to declared component parameters. |
Checkpoints
Each partition's position is persisted as an opaque, adapter-owned token — for Event Hubs, an offset like { "offset": "12345" }. The platform stores and returns these tokens verbatim and never interprets them; only the connection's adapter resolves a position when attaching a reader.
Checkpoints are fenced against operator resets: when you reset a consumer, in-flight processors can no longer overwrite the reset position. See Checkpoint management for reading, resetting, and clearing positions.
Collisions
Stream collisions are per connection+address+subscription.id: only one runtime consumes a subscription at a time. Higher priority wins, then a tenant-created ingress precedes an app resource, then the lexically smaller ingress id. Losers stay installed but inactive with a collision warning signal naming the winner.
Trace propagation
Stream events sent through an egress send carry the sending egress's trace; the ingress restores it so the processing rows connect across the broker. Untrusted partners can opt out per connection with traceContext: { receive: ignore }.
Authorization
ingresses:read # View ingress configurations and checkpoints
ingresses:write # Create, update, and delete ingresses
ingresses:manage # Operational actions: checkpoint resets and clearsRelated
- Checkpoint management — the manage API
- Connections — endpoint and credentials
- Queue ingresses — for work-message queues