Nastavení streamu

Důležité

Tato funkce je ve verzi Public Preview. Správci pracovního prostoru můžou řídit přístup k této funkci ze stránky Previews . Viz Manage Azure Databricks preview.

Stream představuje externí streamový zdroj dat, například Apache Kafka. Streamy ukládají podrobnosti o připojení, ověřování, schémata a konfiguraci příjmu dat. Po vytvoření streamu na něj můžete odkazovat pomocí definic zobrazení funkcí a vytvářet funkce streamování v reálném čase.

Streamy mají názvy tří částí (catalog.schema.stream_name). Přístup ke streamu se řídí přidruženou tabulkou příjmu dat. Podrobnosti viz Příjem a zasypání .

Požadavky

  • Pro spouštění příkazů poznámkového bloku: bezserverový nebo klasický výpočetní cluster se spuštěným modulem Databricks Runtime 17.0 ML nebo novějším.
  • Musí feature-engineering-client být nainstalován balíček Python verze 0.17.0 nebo vyšší.

Připojení ke zdrojům streamu

Před definováním funkcí streamování se připojte a otestujte připojení kanálu Stream Lakeflow k vašemu zprostředkovateli Kafka. Feature Store spoléhá na serverless SDP, což znamená, že budete potřebovat mechanismus pro propojení vašeho klasického výpočetního zařízení (broker nebo endpoint) s Databricks serverless výpočetní technikou. To se děje prostřednictvím produktů jako privatelink nebo tím, že umožňuje přístup ke klasickému výpočetnímu systému z veřejného internetu.

Vytvořte stream

Slouží create_stream() k vytvoření nového streamu. Stream vyžaduje čtyři součásti konfigurace:

  • Konfigurace zdroje: Specifikuje streamovací platformu a specifické detaily zdroje, například předplatné tématu pro zdrojový kód Kafky.
  • Konfigurace připojení: Určuje, jak se připojit a ověřit na platformě streamování, včetně serverů bootstrap a přihlašovacích údajů.
  • Konfigurace schématu: Definuje strukturu klíčů a hodnot zpráv.
  • Konfigurace příjmu dat: Určuje, kde a jak se ingestují data streamu. Podrobnosti viz Příjem a zasypání .

Informace o konfiguraci specifické pro zdroj source_config a nastavení připojení spolu s úplným příkladem create_stream() naleznete v části Apache Kafka. Schéma a možnosti ingestování jsou společné pro všechny zdroje.

Apačský Kafka

Pro streamování z Apache Kafka použijte KafkaStreamConfig jako zdrojovou konfiguraci a pro autentizaci připojení k Unity Catalog. Vizte Streamování ve výpočetním prostředí bez serverů a Připojení k Apache Kafka ohledně připojení ke službě Kafka.

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"
        ),
    ),
)

Režimy předplatného Kafka

Režim odběru určuje, jak Stream vybírá témata Kafka, ze kterých bude číst. Podporují se tři režimy:

Mode Description Příklad
subscribe Čárkami oddělený seznam názvů témat KafkaSubscriptionMode(subscribe="topic1,topic2")
subscribe_pattern Vzor regulárního výrazu v Javě pro porovnávání názvů témat KafkaSubscriptionMode(subscribe_pattern="events-.*")
assign JSON určující přiřazení témat-oddílů KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Kafka autentizace

K autentizaci vůči clusteru Kafka použijte připojení Unity Catalog. Tento přístup se doporučuje pro spravované ověřování. Pokud chcete vytvořit připojení, přečtěte si téma Vytvoření připojení. Autor streamu musí mít u připojení USE CONNECTION. Každý uživatel, který vytváří funkce se Streamem jako zdrojem, musí mít USE CONNECTION také u připojení.

connection_config = StreamConnectionConfig(
    uc_connection_name="my-kafka-connection"
)

Připojení podporuje jak IAM (service credential), tak SASL autentizaci.

IAM (služební průkaz)

Ověřte se například pomocí přihlašovacího oprávnění pro službu Unity Catalog, abyste se mohli připojit k Amazon MSK přes IAM. Pokud chcete vytvořit přihlašovací údaje služby, přečtěte si téma Vytvoření přihlašovacích údajů služby. Nastavte název služebního přihlašovacího oprávnění s credential možností:

CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
    bootstrap_servers '<bootstrap_servers>',
    credential '<service_credential>'
)

Kromě USE CONNECTION u připojení potřebují identity, které používají přihlašovací údaje služby, také ACCESS. Udělení ACCESS odkazu na zmíněnou službu tvůrci Streamu a jakékoli identitě, která se s Streamem objeví. Viz Udělení oprávnění pro přístup k externí cloudové službě pomocí přihlašovacích údajů služby.

SASL

Autentizace SASL používá uživatelské jméno a heslo. Nastavte sasl_mechanism na jednu z následujících možností:

  • PLAIN
  • SCRAM-SHA-256
  • SCRAM-SHA-512

Zadejte přihlašovací údaje pomocí voleb user a password. Připojení tyto přihlašovací údaje bezpečně ukládá.

Následující příklad používá SASL/SCRAM. Pro SASL/PLAIN nastavte sasl_mechanism na 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>'
)

Přímé mTLS

Pro přímé ověřování mTLS zadejte soubory keystore a truststore uložené ve svazku Unity Catalog, přičemž hesla jsou odkazována prostřednictvím oblastí secret scope v Databricks. Další informace o ověřování SSL v systému Kafka najdete v tématu Použití protokolu SSL pro připojení Azure Databricks k Systému 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"
        ),
    ),
)

Konfigurace schématu

Definujte strukturu klíčů a hodnot zpráv tak, aby definice příjmu a funkcí mohly číst jednotlivá pole. U zdrojů Kafka payload_schema odpovídá hodnotě zprávy Kafka (value v modelu klíč–hodnota v systému Kafka) a key_schema odpovídá klíči zprávy Kafka. Musí být zadáno alespoň jedno z payload_schema nebo key_schema.

Každý SchemaConfig přijímá jeden ze tří formátů, které odpovídají tomu, jak zdroj serializuje své zprávy: json_schema, avro_schema, nebo proto_schema. Pokud pro klíč nebo datovou část není k dispozici žádné schéma, považuje se za jednoduchý řetězec.

Příklady kódů v této sekci používají schémata deklarovaná v souladu s DirectSchemas, kde je schéma uvedeno jako řetězec. Pro správu schémat pomocí externího registru schémat viz Schema registry pro podrobnosti.

Schéma JSON

Zadejte řetězec JSON Schema do 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"}'
    ),
)

Schéma Avro

Zadejte řetězec schématu Avro do avro_schema. Podporovány jsou logické typy Avro, včetně timestamp-millis, date, a 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"}}'
            '  ]'
            '}'
        )
    ),
)

Schéma Protobuf

Zadejte ProtoSchemaSpec do proto_schema se zdrojovým textem Protocol Buffers.proto a názvem zprávy datové části. Importovat ProtoSchemaSpec z databricks.feature_engineering.entities.

message_name musí být plně kvalifikované jméno zprávy, včetně deklarovaného package v textu .proto (například com.example.Event, ne Event). Podporovány jsou syntaxe proto2 i proto3.

google.protobuf.Timestamp a jsou podporovány typy skalárních obalů (StringValue, Int32Value, a tak dále) a jejich import je automaticky vyřešen. Jiné známé typy, jako Duration, Struct, a Any, jsou odmítnuty; tyto hodnoty jsou zakódovány jako podporovaný skalár nebo zpráva. Typy fixed32 a fixed64 skalární a map s ne-strunovými klávesami také nejsou podporovány.

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",
        )
    ),
)

Dekódování dat pomocí schémat

Databricks dekóduje každou zprávu pomocí funkcí Spark from_avro, from_protobuf a from_json. Následující chování platí bez ohledu na to, zda deklarujete schéma přímo nebo jej vyřešíte z registru schémat:

  • Deformované záznamy. Dekódování používá režim PERMISSIVE, takže záznam, který neodpovídá svému schématu, se dekóduje jako nulová hodnota místo toho, aby došlo k selhání streamu.
  • Odbory Avro. Sjednocení více typů záznamů dekóduje do struktury s jedním polem na každý typ záznamu, přičemž každé je pojmenováno podle svého Avro záznamu.
  • Protobufové typy. Neznaménková celá čísla se dekódují na širší znaménkový typ (například uint32 na BIGINT a DECIMAL(20,0) na uint64), pole výčtového typu se dekódují na svůj řetězcový název a skalární obalové typy (například StringValue a Int32Value) se dekódují na nulovatelný sloupec obaleného typu.

Registr schématu

Registry schémat uchovávají a verzují schémata, která používají producenti a spotřebitelé streamování, a vynucují pravidla kompatibility, jak se tato schémata vyvíjejí. Když je externí registr schématu konfigurován, Feature Store načte schéma z registru a použije ho k dekódování zprávy o streamování. Při použití registru schématu nedeklarujete schéma přímo ve streamu.

Podpora registru schémat má následující omezení:

  • Podporováno pouze pro Kafka streamy.
  • Podporuje se pouze Confluent Schema Registry
  • Podporovány jsou pouze formáty Avro a Protobuf . Pro čtení JSON zpráv deklarujte schéma přímo v linii. Viz schéma JSON.
  • Každý proud je připojen přesně k jednomu subjektu Confluent pro hodnotu zprávy a jednomu ke klíči zprávy (pokud je uveden). Témata streamů obsahující více záznamů schématu nejsou podporovanou konfigurací. Pokud se váš stream připojuje k tématům, která obsahují více schémat, záznamy, které neodpovídají schématu pro daný subjekt, jsou dekódovány jako null.

Připojte se k registru schématu

Uveďte údaje o připojení k registru jako možnosti v připojení Kafka Unity Catalog a uložte tajný klíč API registru do oboru tajných klíčů Databricks. Identita, pod kterou je Stream spuštěn, musí mít pro rozsah tajných klíčů oprávnění READ, protože ingestní pipeline tento tajný klíč načítá za běhu. Pro způsob, jak vytvořit a nastavit spojení, viz Vytvořit spojení.

Přidejte schema_registry_url, schema_registry_api_key, a schema_registry_api_secret možnosti k připojení používanému pro autentizaci. Následující příklad vytváří Kafka spojení, které se autentizuje s brokerem pomocí přihlašovacího údajů Unity Catalog a s registrem pomocí API klíče:

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>')
)

Nastavte možnost schema_registry_api_secret u připojení Kafka i odkaz na rozsah tajného klíče u streamu na stejný tajný klíč.

Vytvořte stream, který používá registr schématu

Předejte SchemaRegistryConfig jako schema_config. Na tajný klíč rozhraní Registry API odkažte pomocí api_secret_ref a subjekt a formát určete pomocí payload_schema_locator pro hodnotu zprávy nebo pomocí key_schema_locator pro klíč zprávy. Musí být poskytnut alespoň jeden lokalizátor.

Všimněte si zde rozdílů oproti přímým příkladům schémat v sekci konfigurace schémat . Při použití registru schémat neposkytujete schéma přímo v datovém proudu pro schema_config. Místo toho specifikujete schéma SchemaRegistryConfig , které identifikuje schéma v registru.

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"
        ),
    ),
)

Confluentní subjekt je pojmenovaný rozsah, ve kterém je registrována historie verzí schématu a vynucována kompatibilita. Nastavte subject název příslušného rozsahu, který je běžně určen strategií názvů předmětu:

  • TopicNameStrategy (výchozí, odvozuje podmět z názvu tématu): <topic>-value pro hodnotu a <topic>-key pro klíč. Například schéma hodnot pro dané téma transactions používá podmět transactions-value.
  • RecordNameStrategy (odvozuje název subjektu z názvu záznamu ve schématu, nezávisle na topicu): plně kvalifikovaný název záznamu, například com.example.Payment. Toto je jmenný prostor a název záznamu v Avro nebo balíček a název zprávy v Protobuf.
  • TopicRecordNameStrategy (kombinuje názvy témat a záznamů): <topic>-<fully-qualified-record-name>, například transactions-com.example.Payment.

format je povinné. Nastavte to tak, aby SchemaLocatorFormat.FORMAT_AVROSchemaLocatorFormat.FORMAT_PROTOBUF odpovídalo tomu, jak je dané téma serializováno.

Vývoj schématu

Ingestní pipeline při spuštění určí aktuální schéma subjektu. Když v registru schémat zaregistrujete pro daný subjekt novou zpětně kompatibilní verzi schématu, spuštěný kanál zpracování nadále používá verzi, se kterou byl spuštěn.

Protože Databricks spravuje kanál pro příjem dat jako bezserverový kanál Lakeflow, tento kanál se pravidelně restartuje. Při dalším restartu se objeví nová verze schématu. Může trvat až týden, než se v tabulce příjmu objeví nová nebo změněná pole.

O tom, jak pipeline zpracovává záznamy, které neodpovídají aktuálně používanému schématu, viz Dekódování dat pomocí schémat.

Příjem a zasypávání

Parametr ingestion_config určuje, jak se data ze streamu zachycují a ukládají pro trénování a nasazení.

Přístup ke streamu se řídí tabulkou příjmu dat:

  • SELECT v tabulce příjmu dat udělí streamu přístup pro čtení.
  • MANAGE v tabulce příjmu dat uděluje přístup k odstranění.

Další informace o oprávněních k tabulkám najdete v tématech Tabulka a Referenční příručka k oprávněním v katalogu Unity Catalog.

Kanál příjmu dat

Když je stream vytvořen, Databricks spustí řízený ingestion pipeline, který nepřetržitě čte zprávy ze zdrojového proudu a zapisuje je do tabulky Delta (tabulky ingestů). Pipeline začíná z poslední pozice ve zdroji a běží nepřetržitě, přičemž zachycuje pouze nové zprávy, které přicházejí po vytvoření proudu. Tato tabulka příjmu dat se používá pro trénování s funkcemi streamování. Po odstranění datového proudu se odstraní také její kanál příjmu dat a tabulka příjmu dat.

Cíl ingestace

ingestion_destination určuje třídílný název tabulky Delta, do které se zapisují data streamu.

ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
)

Schéma tabulky příjmu dat

Tabulka příjmu obsahuje data zpráv spolu s metadatovými sloupci. Společné sloupce jsou přítomny u každého zdroje; sloupy kafka_* jsou přítomny pouze u Kafkova potoka.

Column Typ Source Description
key Liší se (od key_schema) Common Klíč zprávy, strukturovaný podle schématu, které jste poskytli.
value Liší se (od payload_schema) Common Hodnota zprávy (payload), strukturovaná podle schématu, které jste poskytli.
stream_record_timestamp TIMESTAMP Common Časové razítko záznamu. U dat doplněných dopředu je to časové razítko načtení zdroje. V případě zpětně doplňovaných dat jsou tato data dodaná zákazníkem.
record_source STRING Common Buď "stream" (doplněné dopředu z živého datového proudu), nebo "backfill" (ze zdroje zpětného doplnění dat).
kafka_topic STRING Kafka Téma Kafka, ze které byl záznam spotřebován.
kafka_partition INT Kafka Partice Kafka, ze které byl záznam načten.
kafka_offset LONG Kafka Posun Kafka záznamu v rámci jeho oddílu

Zdroj zpětného doplnění

Protože forward-fill pipeline začíná na nejnovější pozici ve zdroji, nezachytí zprávy, které existovaly před vytvořením datového toku. Chcete-li pro trénink zajistit pokrytí historickými daty, nakonfigurujte volitelný zdroj pro jejich doplnění.

Když je zdroj backfill nakonfigurovaný, Databricks spustí jednorázovou MERGE INTO úlohu, která zkopíruje řádky backfillu do tabulky pro příjem dat s record_source="backfill". Operace MERGE se spustí až poté, co kontrola překrytí potvrdí, že zdroj backfillu a forward-fill stream mají překrývající se časová razítka (viz Překrytí mezi daty backfillu a živého streamu). Pokud není splněna podmínka překrytí do 2 dnů, spustí se funkce MERGE, aby se zabránilo blokování po neomezenou dobu.

Tabulka pro backfill musí obsahovat sloupec stream_record_timestamp typu TIMESTAMP v časovém pásmu UTC. Další sloupce metadat se beze změny přenesou, pokud jsou přítomny ve zdroji zpětného doplnění dat, v opačném případě jsou nastaveny na NULL. Pro Kafka jsou to kafka_topic, kafka_partition a 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"
    ),
)

Překrytí mezi zpětně doplněnými daty a daty z živého streamu

Před spuštěním příkazu MERGE mezi backfill a ingestní tabulkou kontrola překrytí porovnává časová razítka v obou tabulkách:

  • Maximum backfill: Maximální hodnota stream_record_timestamp ve zdroji backfill.
  • Minimum ingestování: Minimální stream_record_timestamp počet řádků (record_source="stream") v ingestní tabulce.

Operace MERGE proběhne, když nejnovější časové razítko backfillu přesáhne nejstarší časové razítko ingestní tabulky alespoň o 1 hodinu. Tím se zajistí, že tabulka příjmu dat neobsahuje žádné mezery. Pokud není splněna podmínka překrytí do 2 dnů, spustí se funkce MERGE, aby se zabránilo blokování po neomezenou dobu.

Protože kanál příjmu dat začíná od poslední pozice ve zdroji, zachycuje pouze zprávy přicházející po vytvoření streamu. Zdroj pro zpětné doplnění dat musí obsahovat data zasahující do časového rozsahu příjmu dat — nejen do času vytvoření datového proudu.

Pokud například vytvoříte datový proud v 15:00, kanál forward-fill začne číst zprávy od 15:00 dále. Zdroj pro backfill musí obsahovat data s časovými razítky minimálně do 16:00 (1 hodinu po začátku forward-fillu), aby kontrola překrytí proběhla úspěšně. To znamená, že byste měli aktualizovat tabulku backfill po 14:00, aby se zajistilo, že tabulka příjmu dat neobsahuje žádné mezery.

Deduplication

Pomocí deduplication_columns určete cesty ke sloupcům pro identifikaci duplicitních řádků při ingesti dat mezi daty backfillu a streamovanými daty forward-fillu. Pro vnořená pole používejte tečkovou notaci (například "value.user_id").

Zvolte sloupce odstranění duplicitních dat na základě vašich dat:

  • Pokud každý záznam ve vašem streamu obsahuje jedinečný identifikátor (například value.transaction_id), použijte tento sloupec pro odstranění duplicitních dat.
  • Pokud zdroj backfill obsahuje kafka_partition a kafka_offset sloupce, použijte je k jednoznačné identifikaci každého záznamu.
  • Pokud nejsou zadány žádné sloupce odstranění duplicit, výchozí klíč odstranění duplicit je úplná kombinace , keyvaluea stream_record_timestamp. Toto se nedoporučuje, protože tato striktní shoda kritérií může snadno vést k duplicitám.
ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
    deduplication_columns=["value.transaction_id"],
)

Správa streamů

Získání datového proudu

stream = client.get_stream(name="my_catalog.my_schema.my_stream")

Seznam streamů

streams = client.list_streams(
    catalog_name="my_catalog",
    schema_name="my_schema",
    max_results=50,
    include_schemas=False,
)

Nastavte include_schemas=True tak, aby zahrnovalo úplné podrobnosti schématu. Schémata můžou být velká a výsledkem může být dlouhotrvající operace. Pokud chcete načíst schémata jednotlivě, použijte get_stream.

Odstranění datového proudu

Odstraněním datového proudu se odstraní také kanál příjmu dat a tabulka příjmu dat.

Warning

Všechny modely nebo funkce, které odkazují na odstraněný stream, už nebudou mít přístup k podkladovým datům datového proudu. Před odstraněním vytvořte kopii tabulky příjmu dat, pokud tato data potřebujete, ale stream už nepotřebujete.

client.delete_stream(name="my_catalog.my_schema.my_stream")

Příklad notebooku

Kompletní příklad, který vytvoří Stream, definuje funkce streamování a nasadí do obslužného koncového bodu, najdete v následujícím poznámkovém bloku:

Notebook s rychlým úvodem do Streaming Feature Views

Pořiďte si notebook