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.
Registering a streaming profile
Section titled “Registering a streaming profile”- Go to Settings → AnyStream and click New AnyStream.
- Choose a broker type. A starter config for that broker appears in the editor.
- Fill in connection details. Sensitive values (passwords, secret keys) should reference an org secret rather than being stored in plaintext — use the
sm:SECRET_NAMEpattern (see Secrets). - Click Test Connection to verify SAMSKARA can reach the broker.
- Save. The profile is now available in every notebook as
sm.stream("profile-name").
Using it in a notebook
Section titled “Using it in a notebook”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.
List topics
Section titled “List topics”topics = client.topics()# ["orders", "payments", "user-events", ...]Read messages
Section titled “Read messages”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 amountorders 2 187431 2026-08-09 07:12 O-123 450.00orders 2 187432 2026-08-09 07:13 O-124 118.50This 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")Read from earliest
Section titled “Read from earliest”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 BronzeWrite messages
Section titled “Write messages”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")Describe a topic
Section titled “Describe a topic”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.
Create a topic
Section titled “Create a topic”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.
Topic defaults and bootstrap
Section titled “Topic defaults and bootstrap”Profile-level defaults and bootstrap topics are configured in Settings → AnyStream when creating or editing a Kafka or Redpanda profile.
Producer and topic defaults
Section titled “Producer and topic defaults”| 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.
Bootstrap topics
Section titled “Bootstrap topics”Topics listed under Bootstrap Topics in Settings are created automatically in two situations:
- On Test Connection — immediately when you click Test Connection while saving the profile.
- 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.
Credentials and secrets
Section titled “Credentials and secrets”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.
Provider config reference
Section titled “Provider config reference”Kafka / Redpanda
Section titled “Kafka / Redpanda”| 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).
Pulsar
Section titled “Pulsar”| 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 |
Kinesis
Section titled “Kinesis”| 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 modes and offset safety
Section titled “Read modes and offset safety”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 toend − limit_per_partition, collects candidates, sorts by_sm_timestampdescending, and returns exactlylimitrows. 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.
offsetis 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.
AI Copilot awareness
Section titled “AI Copilot awareness”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.