# Create a Stream

Launch stage: Beta

`POST /api/2.0/feature-engineering/streams`

Create a Stream, a governed UC entity representing an external streaming data source.

API scopes: mlflow

## Request body

The Stream to create.
- `name` (string, required, ID, Immutable)
  Full three-part (catalog.schema.stream) name of the stream.
- `description` (string, optional)
  User-provided description.
- `source_config` (object, required)
  Source-specific configuration. Determines the streaming platform source.
  - `kafka_stream_config` (object, optional)
    Configuration for Apache Kafka streams.
    - `subscription_mode` (object, required)
      Options to configure which Kafka topics to pull data from.
      - `assign` (string, optional)
        A JSON string that contains the specific topic-partitions to consume from.
         For example, for '{"topicA":[0,1],"topicB":[2,4]}', topicA's 0'th and 1st partitions will be consumed from.
      - `subscribe` (string, optional)
        A comma-separated list of Kafka topics to read from. For example, 'topicA,topicB,topicC'.
      - `subscribe_pattern` (string, optional)
        A regular expression matching topics to subscribe to. For example, 'topic.*' will subscribe to all topics starting with 'topic'.
    - `extra_options` (object, optional)
      Optional Kafka source or consumer options, validated against a server-side
       allowlist at request time. Allowed keys:
       - `maxOffsetsPerTrigger`
       - `startingOffsets`
       - `includeHeaders`
       - `kafka.request.timeout.ms`
       - `kafka.session.timeout.ms`
       - `kafka.max.partition.fetch.bytes`
       The following keys are ingestion-only and are stripped before being forwarded to the materialization pipeline:
       - `maxOffsetsPerTrigger`
       - `startingOffsets`
       Auth and connection details belong on the parent Stream's `connection_config`, not here.
  - `kinesis_stream_config` (object, optional)
    Configuration for AWS Kinesis Data Streams.
    - `stream_names` (object, optional)
      Kinesis stream names to read from.
      - `names` (array of string, optional)
        Kinesis stream names to read from.
    - `stream_arns` (object, optional)
      Kinesis stream ARNs to read from.
      - `arns` (array of string, optional)
        Kinesis stream ARNs to read from. For example,
         'arn:aws:kinesis:us-west-2:111122223333:stream/stream-a'.
    - `extra_options` (object, optional)
      Optional Kinesis source options, validated against a server-side allowlist at request time.
       Allowed keys:
       - `consumerMode`
       - `consumerNamePrefix`
       - `maxFetchRate`
       - `minFetchPeriod`
       - `maxFetchDuration`
       - `maxRecordsPerFetch`
       - `shardsPerTask`
       - `fetchBufferSize`
       - `shardFetchInterval`
       `consumerMode` must be `efo` or `polling` (case-insensitive).
       `maxRecordsPerFetch` applies only during ingestion and does not affect the materialization pipeline.
       Auth and connection details belong on the parent Stream's `connection_config`, not here.
- `connection_config` (object, required)
  Specifies how to connect and authenticate to the stream platform.
  - `uc_connection_name` (string, optional)
    Name of an existing UC Connection for stream platform access.
     Must be the correct type for the streaming platform (e.g. a Kafka Connection for a Kafka
     Stream, or a Kinesis Connection for a Kinesis Stream).
  - `direct_mtls_config` (object, optional)
    Direct mTLS configuration for stream platform access. This is only used in the short term until UC Kafka Connections support mTLS .
     Once UC Kafka Connections support mTLS, this will be deprecated.
    - `bootstrap_servers` (string, required)
      A comma-separated list of host:port pairs for the Kafka bootstrap servers.
    - `mtls_config` (object, required)
      Mutual-TLS authentication configuration.
      - `keystore_location` (string, required)
        Unity Catalog volume path to the JKS keystore file containing the client certificate
         and private key. e.g. "/Volumes/<catalog>/<schema>/<volume>/client.jks". The
         materialization compute must have read permission on this volume.
      - `keystore_password_ref` (object, required)
        Secret-scope reference for the JKS keystore password.
        - `scope` (string, required)
          The Databricks secret scope name.
        - `key` (string, required)
          The key within the scope.
      - `key_password_ref` (object, required)
        Secret-scope reference for the private key password. Often the same value as the
         keystore password (keytool's default), but provided as a separate field because
         Apache Kafka requires it as a distinct option (kafka.ssl.key.password).
        - `scope` (string, required)
          The Databricks secret scope name.
        - `key` (string, required)
          The key within the scope.
      - `truststore_location` (string, required)
        Unity Catalog volume path to the JKS truststore file containing the CA certificate(s)
         trusted to verify the Kafka broker's server certificate.
         e.g. "/Volumes/<catalog>/<schema>/<volume>/truststore.jks".
      - `truststore_password_ref` (object, required)
        Secret-scope reference for the JKS truststore password.
        - `scope` (string, required)
          The Databricks secret scope name.
        - `key` (string, required)
          The key within the scope.
      - `disable_hostname_verification` (boolean, optional)
        Set to true only when the broker certificate's SAN intentionally does not match
         the connection endpoint — for example when reaching the cluster through a
         PrivateLink endpoint whose DNS name is not in the broker certificate. Skipping
         the hostname check removes a defense against man-in-the-middle attacks; do not
         enable casually. mTLS client authentication is unaffected by this option.
        
         See the Apache Kafka SSL security guide for background on this check:
         https://kafka.apache.org/42/security/encryption-and-authentication-using-ssl/#host-name-verification
- `schema_config` (object, required)
  Schema definitions for the stream, provided either directly on the Stream or
   resolved from an external schema registry through a UC Connection.
  - `direct_schemas` (object, optional)
    Schema definitions provided directly on the Stream.
    - `payload_schema` (object, optional)
      Schema for the message payload. For Kafka, this is the value schema.
       Unless the platform supports another schema (e.g. keys for Kafka), this must be specified.
      - `json_schema` (string, optional)
        Schema of the JSON object in standard IETF JSON schema format (https://json-schema.org/).
      - `avro_schema` (string, optional)
        Avro schema in JSON format (https://avro.apache.org/docs/current/specification/).
      - `proto_schema` (object, optional)
        Protocol Buffer schema with its payload message name.
        - `schema_text` (string, required)
          The raw .proto file text (proto2 and proto3 syntax supported, see
           https://protobuf.dev/programming-guides/proto3/ and https://protobuf.dev/programming-guides/proto2/).
        - `message_name` (string, required)
          The fully-qualified name of the message within schema_text that describes the Kafka payload
           (e.g. "Event" or "com.example.Event" if schema_text declares a package). Identifies which
           message is used to decode each Kafka record — a .proto file may declare multiple messages
           but only one represents the payload. Must not be empty.
    - `key_schema` (object, optional)
      Schema for the message key. This is only used for Kafka streams.
       For Kafka, at least one of payload_schema or key_schema must be specified.
      - `json_schema` (string, optional)
        Schema of the JSON object in standard IETF JSON schema format (https://json-schema.org/).
      - `avro_schema` (string, optional)
        Avro schema in JSON format (https://avro.apache.org/docs/current/specification/).
      - `proto_schema` (object, optional)
        Protocol Buffer schema with its payload message name.
        - `schema_text` (string, required)
          The raw .proto file text (proto2 and proto3 syntax supported, see
           https://protobuf.dev/programming-guides/proto3/ and https://protobuf.dev/programming-guides/proto2/).
        - `message_name` (string, required)
          The fully-qualified name of the message within schema_text that describes the Kafka payload
           (e.g. "Event" or "com.example.Event" if schema_text declares a package). Identifies which
           message is used to decode each Kafka record — a .proto file may declare multiple messages
           but only one represents the payload. Must not be empty.
  - `schema_registry_config` (object, optional)
    Resolve schemas from an external schema registry.
    - `uc_connection` (string, optional)
      A Schema Registry UC Connection object.
    - `api_secret_ref` (object, optional)
      Reference to the schema registry API secret in a Databricks secret scope.
       Set this only if required for authentication for the schema registry.
      - `scope` (string, required)
        The Databricks secret scope name.
      - `key` (string, required)
        The key within the scope.
    - `payload_schema_locator` (object, optional)
      Schema locator for the message payload. For Kafka this is the value.
       At least one of payload_schema_locator or key_schema_locator must be set.
      - `confluent_schema` (object, optional)
        Confluent Schema Registry schema locator.
        - `subject` (string, required)
          The Confluent schema registry subject name.
      - `format` (string, required)
        Serialization format for this schema.
        Possible values:
        - `FORMAT_UNSPECIFIED`
        - `FORMAT_AVRO`
        - `FORMAT_PROTOBUF`
        - `FORMAT_JSON`
    - `key_schema_locator` (object, optional)
      Schema locator for the message key. Only used for Kafka streams.
       At least one of payload_schema_locator or key_schema_locator must be set.
      - `confluent_schema` (object, optional)
        Confluent Schema Registry schema locator.
        - `subject` (string, required)
          The Confluent schema registry subject name.
      - `format` (string, required)
        Serialization format for this schema.
        Possible values:
        - `FORMAT_UNSPECIFIED`
        - `FORMAT_AVRO`
        - `FORMAT_PROTOBUF`
        - `FORMAT_JSON`
- `ingestion_config` (object, required)
  Configuration for streaming data ingestion: the managed table storing an offline copy of forward fill data and optional historical backfill.
  - `ingestion_destination` (object, required)
    Destination for the Databricks-managed Delta table that holds an offline copy of the streaming data for querying and training.
     This table contains both 1) forward-filled data from the Stream and 2) backfilled data from the BackfillSource (if provided).
     This table is created and managed by Databricks and is deleted when the Stream is deleted.
    - `delta_table_name` (string, optional)
      The full three-part name (catalog, schema, name) of the Delta table to be created for ingestion.
  - `backfill_source` (object, optional)
    A user-provided source for backfilling data. Historical data is used when creating a training set from streaming features linked to this Stream.
     The backfill data stored in this location will be copied into the ingestion table for offline querying and training.
     The schema for this source must match exactly that of the key and payload schemas specified for this Stream,
     except that it may omit any columns listed in excluded_columns.
    - `delta_table_name` (string, optional)
      The full three-part name (catalog, schema, name) of the Delta table containing the historical data to backfill.
  - `deduplication_columns` (array of string, optional)
    Column paths used to identify duplicate rows during ingestion; only one row per
     distinct combination of these values is kept. Use dot notation for nested fields
     (e.g. `value.user_id`). Empty list means every column is compared.
- `record_type_filter` (string, optional)
  Optional SQL predicate to filter which record types from a streaming channel (e.g. a topic for Kafka) belong to this Stream.
   Events that do not match are not written to the ingestion table and are not used in materialization.
   Example: "value.event_type = 'transaction'".
- `excluded_columns` (array of string, optional)
  Column paths (dot notation, e.g. "value.email" for Kafka) to drop.
   A path may reference a struct, in which case all of its nested fields are dropped (e.g. "value.address" drops "value.address.city" and "value.address.zip").
   These columns are not written to the ingestion table and cannot be referenced by any feature.
   They are dropped from ingestion, backfill, and materialization.
   For direct schemas, each column must exist in the relevant key or payload schema. With a schema registry, a column can be excluded before it exists.
   A column cannot also be a deduplication column in the ingestion_config.

## Returns

Returns the Stream object.

## Response

```json
{
  "name": "string",
  "description": "string",
  "source_config": {
    "kafka_stream_config": {},
    "kinesis_stream_config": {}
  },
  "connection_config": {
    "uc_connection_name": "string",
    "direct_mtls_config": {}
  },
  "schema_config": {
    "direct_schemas": {},
    "schema_registry_config": {}
  },
  "ingestion_config": {
    "ingestion_destination": {},
    "backfill_source": {},
    "deduplication_columns": [
      "string"
    ],
    "ingestion_pipeline_id": "string",
    "ingestion_job_id": 0,
    "backfill_job_id": 0
  },
  "record_type_filter": "string",
  "excluded_columns": [
    "string"
  ],
  "create_time": "string",
  "created_by": "string",
  "update_time": "string",
  "updated_by": "string",
  "browse_only": true
}
```

