Skip to content

AnyStream

AnyStream gives notebooks bounded, read-safe access to streaming brokers — the same credential-management and secret-reference patterns used by database and storage profiles, but aimed at event streams and message queues. Supported brokers: Kafka, Redpanda, Pulsar, NATS, and Kinesis.

“Bounded” means every read() call fetches a fixed number of messages and returns a DataFrame — there is no continuous consumer loop to manage, no offsets to commit, and no risk of touching production consumer group pointers.

  1. Go to Settings → AnyStream and click New AnyStream.
  2. Choose a broker type. A starter config for that broker appears in the editor.
  3. Fill in connection details. Sensitive values (passwords, secret keys) should reference an org secret rather than being stored in plaintext — use the sm:SECRET_NAME pattern (see Secrets).
  4. Click Test Connection to verify SAMSKARA can reach the broker.
  5. Save. The profile is now available in every notebook as sm.stream("profile-name").
client = sm.stream("my-kafka-profile")

sm.stream() returns a StreamClient bound to the named profile. It does not open a connection until you call one of the methods below.

topics = client.topics()
# ["orders", "payments", "user-events", ...]
df = client.read("orders", limit=500)

Returns a DataFrame with one row per message. Metadata columns are prefixed with _sm_ so they never collide with business field names:

Column Description
_sm_topic Topic name
_sm_partition Partition the message came from
_sm_offset Message offset within the partition
_sm_timestamp Broker timestamp in milliseconds
_sm_key Message key (decoded as UTF-8, or null)
_sm_headers Message headers as a dict

When all message values are JSON objects, their fields are expanded as additional top-level columns alongside the _sm_* metadata — no manual json_normalize needed:

_sm_topic _sm_partition _sm_offset _sm_timestamp order_id amount
orders 2 187431 2026-08-09 07:12 O-123 450.00
orders 2 187432 2026-08-09 07:13 O-124 118.50

This makes the ArrowLake bronze pattern natural — Kafka offsets are preserved alongside business data:

df = sm.stream("kafka-prod").read("orders", limit=5000)
sm.write_arrowdelta(df, "bronze.orders_stream")
df = client.read("orders", limit=500, offset="earliest")

Named group for Scheduler-driven ingestion

Section titled “Named group for Scheduler-driven ingestion”

Without group_id, reads use an ephemeral consumer and never commit broker offsets — offset="latest" returns a bounded sample from the current tail of the stream. This is the right mode for interactive inspection and replay.

Supply group_id when you need offset continuity across Scheduler or pipeline runs. Each run continues from where the previous one stopped:

df = client.read("orders", limit=5000, group_id="samskara-orders-bronze")

Named-group reads intentionally advance the group’s committed offset, so they should only be used when that’s the goal. The recommended ingestion template guards against writing an empty result:

client = sm.stream("kafka-prod")
df = client.read("orders", limit=5000, group_id="samskara-orders-bronze")
if not df.empty:
sm.write_arrowdelta(df, "bronze.orders", mode="append")

Scheduled every minute, this gives a near-real-time ingestion pipeline with no additional infrastructure:

Kafka → AnyStream → Scheduler → Notebook → ArrowLake Bronze
result = client.write(df, "processed-orders")
# {"messages_written": 500, "topic": "processed-orders", "duration_ms": 214}

Each row in df is published as a separate message. _sm_* metadata columns are automatically excluded from the message body. To set the Kafka message key from a specific column:

result = client.write(df, "processed-orders", key="order_id")
info = client.describe("orders")
# {
# "name": "orders",
# "provider": "kafka",
# "partition_count": 12,
# "provider_details": { "partitions": [...] }
# }

Returns a provider-neutral envelope so code doesn’t need to branch on broker type. Broker-specific fields are nested under provider_details.

result = client.create_topic("processed-orders")
# {"topic": "processed-orders", "partitions": 1, "replication_factor": 1, "created": True}

Creates the topic if it doesn’t exist, returning {"created": False} if it was already present — safe to call at the top of a notebook. Override the profile defaults inline:

result = client.create_topic(
"high-volume-events",
partitions=12,
replication_factor=3,
retention_ms=86_400_000, # 1 day
)

create_topic() is only available for Kafka and Redpanda profiles. Only parameters passed explicitly override the profile-level defaults.

Profile-level defaults and bootstrap topics are configured in Settings → AnyStream when creating or editing a Kafka or Redpanda profile.

Setting Description Default
Default Partitions Number of partitions for topics created via create_topic() 1
Default Replication Factor Number of replicas 1
Default Retention (days) How long the broker keeps messages 7
Compression Message compression for producers: none, gzip, snappy, lz4, zstd none
Producer Acks Durability guarantee: all (strongest), 1 (leader), 0 (fire-and-forget) all
Auto-create on write Automatically create a topic before the first write() if it doesn’t exist off

Defaults stored on the profile free notebook code from repeating infrastructure parameters — create_topic("my-topic") with no arguments creates a topic with the profile’s configured partition count, replication, and retention.

Topics listed under Bootstrap Topics in Settings are created automatically in two situations:

  1. On Test Connection — immediately when you click Test Connection while saving the profile.
  2. On session start — when a notebook cell calls sm.stream("profile-name") for the first time in a session.

Creation is idempotent: if a topic already exists it is left untouched. Creation failures are non-fatal so a misconfigured broker never prevents a notebook session from starting.

Bootstrap topics are created using the profile’s partition, replication, and retention defaults.

Config values can reference org secrets using the sm:SECRET_NAME pattern:

{
"bootstrap_servers": "broker.example.com:9092",
"security_protocol": "SASL_SSL",
"sasl_mechanism": "SCRAM-SHA-256",
"sasl_username": "sm:KAFKA_USERNAME",
"sasl_password": "sm:KAFKA_PASSWORD"
}

Secrets are resolved at session startup and injected into the runtime environment — the raw secret value is never stored on the profile record.

Field Required Notes
bootstrap_servers Yes host:port, comma-separated for multiple brokers
security_protocol Yes PLAINTEXT, SSL, SASL_PLAINTEXT, or SASL_SSL
sasl_mechanism When SASL PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512
sasl_username When SASL Supports sm: reference
sasl_password When SASL Supports sm: reference
ssl_ca_location When SSL Path to CA certificate on the runner

Redpanda uses the same Kafka protocol — select the Redpanda type for a starter config pre-set for Redpanda’s default SASL port (19092).

Field Required Notes
service_url Yes pulsar://host:6650 or pulsar+ssl://host:6651
token When auth enabled JWT bearer token; supports sm: reference
Field Required Notes
servers Yes List of nats://host:4222 URLs
user When auth enabled Supports sm: reference
password When auth enabled Supports sm: reference
token When token auth Supports sm: reference
Field Required Notes
region Yes AWS region, e.g. us-east-1
aws_access_key_id Yes Supports sm: reference
aws_secret_access_key Yes Supports sm: reference

read() has two distinct modes depending on whether you supply group_id:

Ephemeral (no group_id) — inspection and replay

  • A fresh consumer group (samskara-ephemeral-{uuid}) is created per call and never reused.
  • Broker offsets are never committed. Production consumer groups are untouched.
  • offset="latest" explicitly seeks each partition to end − limit_per_partition, collects candidates, sorts by _sm_timestamp descending, and returns exactly limit rows. You always get the globally most recent messages, not an empty DataFrame.
  • Two sessions reading the same topic simultaneously don’t interfere.

Named group (group_id=) — consumption and pipeline ingestion

  • The named group’s committed offset advances with every call.
  • Each run picks up from where the previous one stopped — this is the correct behaviour for Scheduler-driven ingestion.
  • offset is only used as the reset strategy if the group has no prior committed offset.
  • Use named groups intentionally: sharing a group name with a production consumer will advance its offsets.

The Copilot is aware of your registered AnyStream profiles. When you describe what you want to do with a stream, it will generate code against the correct profile name and broker type — no need to look up config details while writing cells.

  • Secrets — store broker credentials as org secrets and reference them with sm:SECRET_NAME
  • Scheduler — run a notebook that reads from a stream on a schedule
  • ArrowLake — write stream data into a queryable Iceberg table