Note
Access to this page requires authorization. You can try signing in or changing directories.
Access to this page requires authorization. You can try changing directories.
Important
This feature is in Public Preview. Workspace admins can control access to this feature from the Previews page. See Manage Azure Databricks previews.
A Stream represents an external streaming data source, such as Apache Kafka. Streams store connection details, authentication, schemas, and ingestion configuration. After a stream is created, you can reference it using Feature View definitions to create real-time streaming features.
Streams have three-part names (catalog.schema.stream_name). Access to a Stream is governed by its associated ingestion table. See Ingestion and backfill for details.
Requirements
- For running notebook commands: serverless or a classic compute cluster running Databricks Runtime 17.0 ML or above.
- The
feature-engineering-clientPython package version 0.18.0 or above must be installed.
Connecting to stream sources
Before defining streaming features, connect and test a streaming Lakeflow pipeline connection to your Kafka broker. Feature Store relies on serverless SDP which means you will need a mechanism to connect your classic compute (broker or endpoint) to Databricks serverless compute. This is done through products like privatelink or by allowing your classic compute to be accessible from the public internet.
Create a stream
Use create_stream() to create a new Stream. A Stream requires four configuration components:
- Source config: Specifies the streaming platform and source-specific details, such as the topic subscription for a Kafka source.
- Connection config: Specifies how to connect and authenticate to the streaming platform, including bootstrap servers and credentials.
- Schema config: Defines the structure of message keys and values.
- Ingestion config: Specifies where and how stream data is ingested. See Ingestion and backfill for details.
For the source-specific source_config and connection setup, along with a complete create_stream() example, see Apache Kafka. The schema and ingestion options are shared across sources.
Apache Kafka
To stream from Apache Kafka, use KafkaStreamConfig as the source config and a Unity Catalog connection for authentication. See Streaming on serverless compute and Connect to Apache Kafka for Kafka connectivity.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
StreamBackfillSource,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="events-topic"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "transaction_id": {"type": "string"},'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string", "format": "date-time"}'
' }'
'}'
)
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
),
)
Kafka subscription modes
The subscription mode specifies how the Stream selects Kafka topics to consume from. Three modes are supported:
| Mode | Description | Example |
|---|---|---|
subscribe |
Comma-separated list of topic names | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Java regex pattern matching topic names | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
JSON specifying topic-partition assignments | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Kafka authentication
Unity Catalog connection (recommended)
Use a Unity Catalog connection to authenticate to your Kafka cluster. This is the recommended approach for managed authentication. To create a connection, see Create a connection. The Stream's creator must have USE CONNECTION on the connection. Any user materializing features with the Stream as a source must also have USE CONNECTION on the connection.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
The connection supports both IAM (service credential) and SASL authentication.
IAM (service credential)
Authenticate with a Unity Catalog service credential, for example to connect to Amazon MSK with IAM. To create a service credential, see Create service credentials. Set the service credential name with the credential option:
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>'
)
In addition to USE CONNECTION on the connection, identities that use the service credential need ACCESS on it. Grant ACCESS on the referenced service credential to the Stream's creator and to any identity that materializes features with the Stream. See Grant permissions to use a service credential to access an external cloud service.
SASL
SASL authentication uses a username and password. Set sasl_mechanism to one of the following:
PLAINSCRAM-SHA-256SCRAM-SHA-512
Provide the credentials with the user and password options. The connection stores these credentials securely.
The following example uses SASL/SCRAM. For SASL/PLAIN, set sasl_mechanism to PLAIN.
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
sasl_mechanism 'SCRAM-SHA-512',
user '<username>',
password '<password>'
)
Direct mTLS
For direct mTLS authentication, provide keystore and truststore files stored on a Unity Catalog volume, with passwords referenced through Databricks secret scopes. For more information on SSL authentication with Kafka, see Use SSL to connect Azure Databricks to Kafka.
from databricks.feature_engineering.entities import (
DirectMtlsConfig,
MtlsConfig,
SecretScopeReference,
)
connection_config = DirectMtlsConfig(
bootstrap_servers="broker1:9092,broker2:9092",
mtls_config=MtlsConfig(
keystore_location="/Volumes/my_catalog/my_schema/my_volume/keystore.jks",
keystore_password_ref=SecretScopeReference(
scope="my_scope", key="keystore_password"
),
key_password_ref=SecretScopeReference(
scope="my_scope", key="key_password"
),
truststore_location="/Volumes/my_catalog/my_schema/my_volume/truststore.jks",
truststore_password_ref=SecretScopeReference(
scope="my_scope", key="truststore_password"
),
),
)
Schema configuration
Define the structure of message keys and values so that ingestion and feature definitions can read individual fields. For Kafka sources, payload_schema corresponds to the Kafka message value (the value in Kafka's key-value model) and key_schema corresponds to the Kafka message key. At least one of payload_schema or key_schema must be provided.
Each SchemaConfig accepts one of three formats, matching how the source serializes its messages: json_schema, avro_schema, or proto_schema. If no schema is provided for a key or payload, it is treated as a simple string.
The code examples in this section use schemas declared inline with DirectSchemas, where the schema is provided as a string. To manage schemas using an external schema registry see Schema registry for details.
JSON schema
Provide a JSON Schema string to json_schema.
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"}'
' }'
'}'
)
),
key_schema=SchemaConfig(
json_schema='{"type": "string"}'
),
)
Avro schema
Provide an Avro schema string to avro_schema. Avro logical types are supported, including timestamp-millis, date, and decimal.
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
avro_schema=(
'{'
' "type": "record",'
' "name": "Event",'
' "fields": ['
' {"name": "user_id", "type": "string"},'
' {"name": "amount", "type": "double"},'
' {"name": "event_time",'
' "type": {"type": "long", "logicalType": "timestamp-millis"}}'
' ]'
'}'
)
),
)
Protobuf schema
Provide a ProtoSchemaSpec to proto_schema with the Protocol Buffers .proto source text and the payload message name. Import ProtoSchemaSpec from databricks.feature_engineering.entities.
message_name must be the fully-qualified message name, including the package declared in the .proto text (for example, com.example.Event, not Event). Both proto2 and proto3 syntax are supported.
google.protobuf.Timestamp and the scalar wrapper types (StringValue, Int32Value, and so on) are supported, and their imports are resolved automatically. Other well-known types, such as Duration, Struct, and Any, are rejected; encode those values as a supported scalar or message instead. The fixed32 and fixed64 scalar types and map with non-string keys are also not supported.
from databricks.feature_engineering.entities import ProtoSchemaSpec
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
proto_schema=ProtoSchemaSpec(
schema_text=(
'syntax = "proto3";\n'
'package com.example;\n'
'import "google/protobuf/timestamp.proto";\n'
'message Event {\n'
' string user_id = 1;\n'
' double amount = 2;\n'
' google.protobuf.Timestamp event_time = 3;\n'
'}'
),
message_name="com.example.Event",
)
),
)
Decoding data using schemas
Databricks decodes each message with Spark's from_json, from_avro, and from_protobuf functions. The following behaviors apply whether you declare the schema inline or resolve it from a schema registry:
- Malformed records. Decoding uses the
PERMISSIVEmode, so a record that does not match its schema decodes to a null value instead of failing the stream. - Avro unions. A union of multiple record types decodes to a struct with one field per record type, each named after its Avro record.
- Protobuf types. Unsigned integers decode to a wider signed type (for example,
uint32toBIGINTanduint64toDECIMAL(20,0)), enum fields decode to their string name, and scalar wrapper types (for example,StringValueandInt32Value) decode to a nullable column of the wrapped type.
Schema registry
Schema registries store and version schemas that streaming producers and consumers use, enforcing compatibility rules as those schemas evolve. When an external schema registry is configured, Feature Store reads the schema from the registry and uses it to decode the streaming message. You do not declare the schema inline on the Stream when using a schema registry.
Schema registry support has the following limitations:
- Supported for Kafka streams only.
- Only Confluent Schema Registry is supported
- Only the Avro and Protobuf formats are supported. To read JSON messages, declare the schema inline instead. See JSON schema.
- Each Stream is connected to exactly one Confluent subject for the message value, and one for the message key (if provided). Stream topics containing multiple schema records is not a supported configuration. If your Stream connects to topics that contain multiple schemas, records that do not match the schema for the specified subject are decoded as null.
Connect to a schema registry
Provide the registry connection details as options on the Kafka Unity Catalog connection, and store the registry API secret in a Databricks secret scope. The Stream's run-as identity must have READ permission on the secret scope, because the ingestion pipeline reads the secret at runtime. For how to create and configure a connection, see Create a connection.
Add the schema_registry_url, schema_registry_api_key, and schema_registry_api_secret options to the connection used for authentication. The following example creates a Kafka connection that authenticates to the broker with a Unity Catalog service credential and to the registry with an API key:
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>',
schema_registry_url 'https://<registry-host>',
schema_registry_api_key '<registry_api_key>',
schema_registry_api_secret secret('<scope>', '<key>')
)
Set both the schema_registry_api_secret option on the Kafka connection and the secret scope reference on the Stream to the same secret.
Create a stream that uses a schema registry
Pass a SchemaRegistryConfig as the schema_config. Reference the registry API secret with api_secret_ref, and identify the subject and format with payload_schema_locator for the message value, or key_schema_locator for the message key. At least one locator must be provided.
Note the differences here compared to the direct schema examples in the Schema configuration section. When using a schema registry, you do not provide the schema inline on the Stream to schema_config. Instead, you specify a SchemaRegistryConfig that identifies the schema in the registry.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
SchemaRegistryConfig,
SchemaLocator,
SchemaLocatorConfluentSchema,
SchemaLocatorFormat,
SecretScopeReference,
IngestionConfig,
IngestionDestination,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="transactions"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=SchemaRegistryConfig(
api_secret_ref=SecretScopeReference(
scope="my_scope", key="sr_api_secret"
),
payload_schema_locator=SchemaLocator(
confluent_schema=SchemaLocatorConfluentSchema(
subject="transactions-value"
),
format=SchemaLocatorFormat.FORMAT_AVRO,
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.transactions_ingestion"
),
),
)
A Confluent subject is the named scope under which a schema's version history is registered and compatibility is enforced. Set subject to the relevant scope's name, which is commonly determined from the subject name strategy:
- TopicNameStrategy (default, derives the subject from the topic name):
<topic>-valuefor the value and<topic>-keyfor the key. For example, the value schema for the topictransactionsuses the subjecttransactions-value. - RecordNameStrategy (derives the subject from the schema's record name, independent of the topic): the fully-qualified record name, such as
com.example.Payment. This is the record's namespace and name for Avro, or the message's package and name for Protobuf. - TopicRecordNameStrategy (combines the topic and record names):
<topic>-<fully-qualified-record-name>, such astransactions-com.example.Payment.
format is required. Set it to SchemaLocatorFormat.FORMAT_AVRO or SchemaLocatorFormat.FORMAT_PROTOBUF to match how the topic is serialized.
Schema evolution
The ingestion pipeline resolves the subject's current schema when it starts. When you register a new backward-compatible schema version on the subject in the schema registry, the running pipeline continues to use the version it started with.
For schema registry-backed Streams, the ingestion pipeline restarts automatically every couple of hours. On each restart it picks up the subject's latest schema version, and new or changed fields appear in the ingestion table.
Streams that use direct schemas instead of a schema registry evolve their schema with update_stream. See Update a stream.
For how the pipeline handles records that don't match the schema it is currently using, see Decoding data using schemas.
Filter records by type
A Stream decodes each record using a single key and value schema (if provided), whether you specify them directly or use a schema registry. Because a topic can carry more than one type of record and Streams can subscribe to multiple topics, use record_type_filter to select which records from the topic belong to this Stream.
Provide a SQL expression that references decoded fields with dot notation, for example value.event_type = 'transaction'. Records that don't match the filter are ignored. They aren't written to the ingestion table and aren't used in materialization. To create a Stream for other record types, create a separate Stream with a different record_type_filter.
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
record_type_filter="value.event_type = 'transaction'",
)
Even without record_type_filter, decoding never fails the Stream. A record that doesn't match the configured schema is decoded permissively. The records decode in one of the following ways:
- Into a row with
NULLvalues for the fields that the schema expects but the record omits (JSON, Avro, and Protobuf). - Into a row holding values that belong to a different record type (Avro and Protobuf only).
To identify whether a row in the ingestion table belongs to the expected record type, use one of the following checks:
- Check that a field equals an expected value, for example
value.event_type = 'transaction'(preferred for Avro and Protobuf). - Check that a field is non-
NULL, for examplevalue.activity_id IS NOT NULL.
Using record_type_filter with separate Streams is recommended when the schemas differ substantially between record types on the topic or you want to govern access to each record type independently. To maintain costs, Databricks recommends that you keep a small number of Streams, since each Stream has a separate ingestion pipeline and ingestion table. Each Stream also uses separate compute at materialization time. You can use feature-specific filters for materialization.
record_type_filter is different from a feature's filter_condition. record_type_filter is set on the Stream and controls which records are ingested and available to all features using the Stream as a source, whereas filter_condition is set on an individual feature and filters rows before aggregation. See Filter conditions on streaming sources for more details on filter_condition.
Ingestion and backfill
The ingestion_config parameter configures how stream data is captured and stored for training and serving.
Access to a Stream is governed by the ingestion table:
SELECTon the ingestion table grants read access to the Stream.MANAGEon the ingestion table grants delete access.
For more information on table privileges, see Table and Unity Catalog privileges reference.
Ingestion pipeline
When a stream is created, Databricks starts a managed ingestion pipeline that continuously reads messages from the source stream and writes them into a Delta table (the ingestion table). The pipeline starts from the latest position in the source and runs continuously, capturing only new messages that arrive after the stream is created. This ingestion table is used for training with streaming features. When a stream is deleted, its ingestion pipeline and ingestion table are also deleted.
Ingestion destination
The ingestion_destination specifies the three-part Delta table name where stream data is written.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Ingestion table schema
The ingestion table contains the message data along with metadata columns. The common columns are present for every source; the kafka_* columns are present only for a Kafka stream.
| Column | Type | Source | Description |
|---|---|---|---|
key |
Varies (from key_schema) |
Common | The message key, structured according to the schema you provided. |
value |
Varies (from payload_schema) |
Common | The message value (payload), structured according to the schema you provided. |
stream_record_timestamp |
TIMESTAMP |
Common | The record timestamp. For forward-fill data, this is the source ingest timestamp. For backfill data, this is customer-supplied. |
record_source |
STRING |
Common | Either "stream" (forward-fill from the live stream) or "backfill" (from the backfill source). |
kafka_topic |
STRING |
Kafka | The Kafka topic the record was consumed from. |
kafka_partition |
INT |
Kafka | The Kafka partition the record was consumed from. |
kafka_offset |
LONG |
Kafka | The Kafka offset of the record within its partition. |
Backfill source
Because the forward-fill pipeline starts from the latest position in the source, it does not capture messages that existed before the stream was created. To provide historical data coverage for training, configure an optional backfill source.
When a backfill source is configured, Databricks runs a one-time MERGE INTO job that copies backfill rows into the ingestion table with record_source="backfill". The MERGE runs only after the overlap checker confirms that the backfill source and the forward-fill stream have overlapping timestamps (see Overlap between backfill and live stream data). If the overlap condition is not met within 2 days, the MERGE runs anyway to avoid blocking indefinitely.
The backfill table must include a stream_record_timestamp column of type TIMESTAMP in UTC timezone. Other metadata columns are passed through if present on the backfill source, or set to NULL otherwise. For Kafka, these are kafka_topic, kafka_partition, and kafka_offset.
from databricks.feature_engineering.entities import StreamBackfillSource
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
backfill_source=StreamBackfillSource(
delta_table_name="my_catalog.my_schema.historical_events"
),
)
Overlap between backfill and live stream data
Before running a MERGE between backfill and the ingestion table, an overlap check compares timestamps on the two tables:
- Backfill max: The maximum
stream_record_timestampin the backfill source. - Ingestion min: The minimum
stream_record_timestampof rows (record_source="stream") in the ingestion table.
The MERGE proceeds when the backfill's latest timestamp exceeds the ingestion table's earliest timestamp by at least 1 hour. This overlap ensures there are no gaps in the ingestion table. If the overlap condition is not met within 2 days, the MERGE runs anyway to avoid blocking indefinitely.
Because the ingestion pipeline starts from the latest position in the source, it only captures messages arriving after the stream is created. Your backfill source must contain data that extends into the ingestion time range — not just up to the stream creation time.
For example, if you create a stream at 3:00 PM, the forward-fill pipeline begins reading messages from 3:00 PM onward. Your backfill source must include data with timestamps through at least 4:00 PM (1 hour past the forward-fill start) to satisfy the overlap check. This means you should update your backfill table after 4:00 pm to ensure ingestion table has no gaps.
Deduplication
Use deduplication_columns to specify column paths for identifying duplicate rows during ingestion between backfill and forward-fill stream data. Use dot notation for nested fields (for example, "value.user_id").
Choose deduplication columns based on your data:
- If each record in your stream contains a unique identifier (for example,
value.transaction_id), use that column for deduplication. - If your backfill source includes
kafka_partitionandkafka_offsetcolumns, use those to uniquely identify each record. - If no deduplication columns are specified, the default deduplication key is the full combination of
key,value, andstream_record_timestamp. This is not recommended as this strict criteria matching can easily lead to duplicates.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Cost attribution
Set tags and budget_policy_id on the IngestionConfig to attribute the cost of the Stream's managed ingestion. Azure Databricks applies them to the ingestion Lakeflow pipeline and its forward-fill and backfill jobs when the Stream is created.
For an example, the tag limits, and how to query the attributed spend, see Attribute costs with tags and serverless usage policies.
Exclude columns from a Stream
Use excluded_columns to drop specific columns from a Stream that you don't want to ingest. An excluded column is not written to the ingestion table and cannot be referenced by a feature or used in training.
Specify each column using dot notation into the message key or value, such as value.user.email or key.account_id. These columns are dropped from the decoded key and value across ingestion, backfill, and materialization. If a path points to a struct, all of its nested fields are dropped as well (for example, value.address also drops value.address.city and value.address.zip).
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
excluded_columns=["value.user.email", "value.user.ssn"],
)
When using direct schemas, the excluded column must already exist in the key or value schema, or create_stream fails. When using a schema registry, you can exclude a column before it exists. An excluded column also can't be a deduplication column, because deduplication columns are required to identify duplicate rows. Any feature that references an excluded column (for example, as an entity, timeseries, or input) fails to be created.
You can change a Stream's excluded columns after creation with update_stream, on both direct schema and schema registry-backed Streams. See Update a stream for more details.
Manage streams
Get a stream
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
List streams
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
Set include_schemas=True to include full schema details. Schemas can be large and this might result in a long-running operation. To retrieve schemas individually instead, use get_stream.
Update a stream
Use update_stream to change a Stream after creation. Pass schema_config to evolve a direct schema, excluded_columns to change which columns are dropped, or both. Updating other fields is not supported. Create a new Stream instead.
Updating a Stream restarts its ingestion pipeline so the change takes effect. Ingestion typically resumes within a few minutes.
Evolve a direct schema
For a Stream that uses direct schemas, pass a DirectSchemas to schema_config. Set payload_schema, key_schema, or both. A side you don't set is left unchanged. Schema registry-backed Streams reject a schema_config update and must be evolved through the registry instead.
from databricks.feature_engineering.entities import DirectSchemas, SchemaConfig
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"},'
' "channel": {"type": "string"}'
' }'
'}'
)
),
),
)
Schema updates must be backward-compatible so that the running ingestion pipeline can keep decoding existing records and writing to the ingestion table. Any other changes are rejected.
What's allowed depends on the format:
- JSON and Protobuf: add optional fields, remove fields, and widen a field's type (for example,
inttobigint). Protobuf also allows reordering fields. - Avro: allows only widening
inttolongand removing a trailing field whose bytes no later field reads. To evolve an Avro schema more freely, use a schema registry-backed stream instead.
Adding fields grows the ingestion table's decoded key and value structs. Rows written before the update keep their original shape, and the added fields read as NULL for those earlier rows. Removals and type changes take effect only for records ingested after the update.
Change excluded columns
Pass the complete new set of column paths to excluded_columns, which replaces the existing set. Pass an empty list ([]) to clear all exclusions. For details on this behavior, see Exclude columns from a Stream.
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
excluded_columns=["value.user.email", "value.user.ssn"],
)
Changing excluded columns is forward-only. Newly excluded columns stop being written (appearing at NULL) and newly included ones start being populated going forward, while previously written rows are left as-is. To prevent a new column from ever being ingested:
- Schema registry: add the column to
excluded_columnsfirst and wait for the ingestion pipeline to restart, then register the new schema version in the registry. - Direct schemas: add the column to
schema_configand toexcluded_columnsin the sameupdate_streamcall.
Delete a stream
Deleting a stream also deletes its ingestion pipeline and ingestion table.
Warning
Any models or features that reference the deleted stream will no longer have access to the underlying stream data. Create a copy of the ingestion table before deletion if you need this data but no longer need the stream.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Example notebook
For an end-to-end example that creates a Stream, defines streaming features, and deploys to a serving endpoint, see the following notebook: