Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Ważne
Ta funkcja jest dostępna w publicznej wersji testowej. Administratorzy obszaru roboczego mogą kontrolować dostęp do tej funkcji ze strony Podglądy . Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.
Stream oznacza zewnętrzne źródło danych strumieniowych, takie jak Apache Kafka. Strumienie przechowują parametry połączenia, dane uwierzytelniające, schematy i konfigurację pozyskiwania danych. Po utworzeniu strumienia można odwoływać się do niego przy użyciu definicji widoku funkcji w celu utworzenia funkcji przesyłania strumieniowego w czasie rzeczywistym.
Strumienie mają trzyczęściowe nazwy (catalog.schema.stream_name). Dostęp do strumienia jest kontrolowany przez skojarzoną z nim tabelę pozyskiwania danych. Aby uzyskać szczegółowe informacje, zobacz Pozyskiwanie danych i uzupełnianie danych.
Requirements
- Uruchamianie poleceń notesu: bezserwerowy lub klasyczny klaster obliczeniowy z uruchomionym środowiskiem Databricks Runtime 17.0 ML lub nowszym.
- Pakiet
feature-engineering-clientPython w wersji 0.17.0 lub wyższej musi być zainstalowany.
Tworzenie strumienia
Użyj create_stream(), aby utworzyć nowy Stream. Usługa Stream wymaga czterech składników konfiguracji:
- Konfiguracja źródła: określa platformę przesyłania strumieniowego (na przykład Kafka) i szczegóły specyficzne dla źródła (takie jak subskrypcja tematu dla platformy Kafka).
- Konfiguracja połączenia: określa sposób nawiązywania połączenia i uwierzytelniania z platformą przesyłania strumieniowego, w tym serwerów bootstrap i poświadczeń.
- Konfiguracja schematu: definiuje strukturę kluczy i wartości komunikatów.
- Konfiguracja pozyskiwania danych: określa, gdzie i w jaki sposób pozyskiwane są dane strumieniowe. Aby uzyskać szczegółowe informacje, zobacz Pozyskiwanie danych i uzupełnianie danych.
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"
),
),
)
Łączenie ze źródłami strumieniowymi
Przed zdefiniowaniem funkcji przesyłania strumieniowego połącz się i przetestuj połączenie potoku przesyłania strumieniowego Lakeflow z brokerem platformy Kafka. Zobacz Przesyłanie strumieniowe w środowisku bezserwerowym i Łączenie z platformą Apache Kafka.
Informacje o usłudze AWS Managed Streaming (Amazon MSK) można znaleźć w dokumencie Bezserwerowa prywatna łączność z usługą Amazon MSK. Aby uzyskać szczegółowe informacje na temat opcji uwierzytelniania platformy Kafka, zobacz Uwierzytelnianie.
Authentication
Połączenie z usługą Unity Catalog (zalecane)
Użyj połączenia Unity Catalog, aby uwierzytelnić się w klastrze Kafka. Jest to zalecane podejście do uwierzytelniania zarządzanego. Aby utworzyć połączenie, zobacz Tworzenie połączenia. Twórca strumienia musi mieć uprawnienie USE CONNECTION do połączenia. Każdy użytkownik, który realizuje funkcje z Streamem jako źródłem, musi również mieć USE CONNECTION na połączeniu.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
Połączenie obsługuje zarówno uwierzytelnianie IAM (service credential), jak i SASL.
IAM (poświadczenie służbowe)
Uwierzytelnij się przy użyciu poświadczeń usługi Unity Catalog, na przykład aby połączyć się z usługą Amazon MSK za pomocą mechanizmu IAM. Aby utworzyć poświadczenia usługi, zobacz Tworzenie poświadczeń usługi. Ustaw nazwę poświadczenia usługi z opcją credential :
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>'
)
Oprócz USE CONNECTION dla połączenia tożsamości korzystające z poświadczenia usługi potrzebują również ACCESS na tym poświadczeniu. Nadaj ACCESS do wskazanego poświadczenia usługi twórcy strumienia oraz każdej tożsamości, która materializuje funkcje za pomocą strumienia. Zobacz Przyznawanie uprawnień do używania poświadczenia usługi w celu uzyskania dostępu do zewnętrznej usługi w chmurze.
SASL
Uwierzytelnianie SASL wykorzystuje nazwę użytkownika i hasło. Ustaw sasl_mechanism na jedną z następujących wartości:
PLAINSCRAM-SHA-256SCRAM-SHA-512
Podaj poświadczenia za pomocą opcji user i password. Połączenie bezpiecznie przechowuje te dane uwierzytelniające.
Poniższy przykład wykorzystuje SASL/SCRAM. Dla SASL/PLAIN ustaw 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>'
)
Bezpośredni mTLS
W przypadku bezpośredniego uwierzytelniania mTLS podaj pliki keystore i truststore przechowywane w woluminie Unity Catalog, z hasłami, do których odwołują się zakresy wpisów tajnych Databricks. Aby uzyskać więcej informacji na temat uwierzytelniania SSL na platformie Kafka, zobacz Używanie protokołu SSL do nawiązywania połączenia Azure Databricks z platformą 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"
),
),
)
Tryby subskrypcji
Tryb subskrypcji określa sposób wybierania tematów platformy Kafka do korzystania z usługi Stream. Obsługiwane są trzy tryby:
| Tryb | Description | Przykład |
|---|---|---|
subscribe |
Rozdzielona przecinkami lista nazw tematów | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Nazwy tematów pasujące do wyrażenia regularnego Java | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
Kod JSON określający przypisania partycji tematu | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Konfiguracja schematu
Zdefiniuj strukturę kluczy i wartości wiadomości, aby definicje odbiorców i cech mogły odczytywać poszczególne pola. W przypadku źródeł payload_schema platformy Kafka odpowiada wartości komunikatu platformy Kafka ( value w modelu klucz-wartość platformy Kafka) i key_schema odpowiada kluczowi komunikatu platformy Kafka. Należy podać co najmniej jedno z pól payload_schema lub key_schema.
Każdy SchemaConfig z nich akceptuje jeden z trzech formatów, odpowiadających sposobowi serializacji wiadomości json_schemaprzez źródło: , avro_schema, lub proto_schema. Jeśli dla klucza lub ładunku nie podano żadnego schematu, jest on traktowany jako prosty ciąg.
Przykłady kodu w tej sekcji używają schematów zadeklarowanych bezpośrednio przy użyciu DirectSchemas, w których schemat jest podany jako ciąg znaków. Aby zarządzać schematami za pomocą zewnętrznego rejestru schematów, szczegóły można znaleźć w rejestrze schematów .
Schemat systemu JSON
Dostarcz ciąg schematu JSON 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"}'
),
)
Schemat Avro
Dostarcz ciąg schematu Avro do avro_schema. Obsługiwane są typy logiczne Avro, w tym timestamp-millis, date, oraz 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"}}'
' ]'
'}'
)
),
)
Schemat Protobuf
Przekaż do ProtoSchemaSpecproto_schema tekst źródłowy Protocol Buffers.proto oraz nazwę komunikatu payload. Zaimportuj ProtoSchemaSpec z databricks.feature_engineering.entitiespliku .
message_name musi być w pełni kwalifikowaną nazwą komunikatu, wraz z package zadeklarowaną w tekście .proto (na przykład com.example.Event, a nie Event). Obsługiwane są zarówno składnie proto2, jak i proto3.
Obsługiwane są google.protobuf.Timestamp oraz typy opakowujące dla typów skalarnych (StringValue, Int32Value itd.), a ich importy są rozwiązywane automatycznie. Inne znane typy, takie jak Duration, Struct, i Any, są odrzucane; koduj te wartości jako wspierany skalar lub komunikat. Typy skalarne fixed32 i fixed64 oraz map z kluczami innymi niż ciągi znaków również nie są obsługiwane.
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",
)
),
)
Dekodowanie danych za pomocą schematów
Databricks dekoduje każdą wiadomość za pomocą funkcji Spark from_avro, from_protobuf i from_json. Następujące zachowania mają zastosowanie niezależnie od tego, czy deklarujesz schemat w linii, czy rozwiązujesz go z rejestru schematów:
- Zdeformowane rekordy. Dekodowanie używa trybu
PERMISSIVE, więc rekord, który nie jest zgodny ze swoim schematem, jest dekodowany do wartości null zamiast powodować błąd strumienia. - Związki zawodowe Avro. Unia kilku typów rekordów jest dekodowana do struktury z jednym polem dla każdego typu rekordu, z których każde nosi nazwę odpowiadającą rekordowi Avro.
- Typy protobufów. Liczby całkowite bez znaku są dekodowane do szerszego typu ze znakiem (na przykład
uint32doBIGINTorazuint64doDECIMAL(20,0)), pola enum są dekodowane do ich nazw tekstowych, a skalarne typy opakowujące (na przykładStringValueiInt32Value) są dekodowane do kolumny dopuszczającej wartość null typu opakowanego.
Rejestr schematów
Rejestry schematów przechowują i wersjonują schematy używane przez producentów i konsumentów danych strumieniowych, wymuszając reguły zgodności w miarę ich ewolucji. Gdy zewnętrzny rejestr schematu zostanie skonfigurowany, Feature Store odczytuje schemat z rejestru i używa go do dekodowania wiadomości streamingowej. Nie deklarujesz schematu w linii w Streamie, korzystając z rejestru schematów.
Wsparcie rejestru schematów ma następujące ograniczenia:
- Obsługiwany jest tylko Confluent Schema Registry
- Obsługiwane są tylko formaty Avro i Protobuf . Aby czytać komunikaty JSON, zadeklaruj schemat w linii. Zobacz schemat JSON.
- Każdy strumień jest połączony dokładnie z jednym subjectem Confluent dla wartości komunikatu oraz jednym dla klucza komunikatu (jeśli został podany). Tematy strumienia zawierające wiele rekordów schematu nie stanowią obsługiwanej konfiguracji. Jeśli Twój Strumień łączy się z tematami zawierającymi wiele schematów, rekordy niezgodne ze schematem dla danego tematu są dekodowane jako null.
Połącz się z rejestrem schematów
Udostępnij szczegóły połączenia rejestru jako opcje na połączeniu Kafka Unity Catalog i przechowuj sekret API rejestru w skali tajnego Databricks. Tożsamość uruchomieniowa strumienia musi mieć READ uprawnienie do zakresu wpisów tajnych, ponieważ potok pozyskiwania danych odczytuje wpis tajny w czasie wykonywania. Aby dowiedzieć się, jak utworzyć i skonfigurować połączenie, zobacz: Utwórz połączenie.
Dodaj schema_registry_url, schema_registry_api_key, i schema_registry_api_secret opcje do połączenia używanego do uwierzytelniania. Poniższy przykład tworzy połączenie Kafki, które uwierzytelnia się do brokera za pomocą uwierzytelniania usługi Unity Catalog oraz do rejestru za pomocą klucza API:
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>')
)
Ustaw schema_registry_api_secret zarówno opcję połączenia Kafka, jak i referencję do zakresu tajnego w strumieniu na ten sam sekret.
Stwórz strumień korzystający z rejestru schematu
Przekaż SchemaRegistryConfig jako schema_config. Odwołaj się do sekretu interfejsu API rejestru za pomocą api_secret_ref, a temat i format określ za pomocą key_schema_locator dla wartości wiadomości lub payload_schema_locator dla klucza wiadomości. Należy podać co najmniej jeden lokalizator.
Zwróć uwagę na różnice w porównaniu z przykładami schematów bezpośrednich w sekcji konfiguracji schematów . Korzystając z rejestru schematów, nie udostępniasz schematu w linii na strumieniu do schema_config. Zamiast tego podajesz SchemaRegistryConfig, który identyfikuje schemat w rejestrze.
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"
),
),
)
Podmiot Confluent to nazwany zakres, w którym rejestrowana jest historia wersji schematu i wymuszona jest zgodność. Ustaw subject na nazwę odpowiedniego zakresu, która jest zwykle określana na podstawie strategii nazewnictwa tematu:
-
TopicNameStrategy (domyślnie, wyprowadza podmiot z nazwy tematu):
<topic>-valuedla wartości i<topic>-keydla klucza. Na przykład schemat wartości dla tematutransactionswykorzystuje podmiottransactions-value. -
RecordNameStrategy (wyprowadza nazwę subjectu z nazwy rekordu schematu, niezależnie od topicu): w pełni kwalifikowana nazwa rekordu, na przykład
com.example.Payment. Jest to przestrzeń nazw i nazwa rekordu w Avro lub pakiet i nazwa wiadomości w Protobuf. -
TopicRecordNameStrategy (łączy nazwy tematów i rekordów):
<topic>-<fully-qualified-record-name>, takie jaktransactions-com.example.Payment.
Ciąg format jest wymagany. Ustaw tę wartość na SchemaLocatorFormat.FORMAT_AVRO lub SchemaLocatorFormat.FORMAT_PROTOBUF, aby odpowiadała sposobowi serializacji tematu.
Ewolucja schematu
Potok przetwarzania danych wejściowych ustala bieżący schemat obiektu w momencie uruchomienia. Gdy rejestrujesz nową, kompatybilną wsteczną wersję schematu na danym obszarze w rejestrze schematów, działający potok nadal korzysta z wersji, od której zaczął.
Ponieważ Databricks zarządza potokiem pozyskiwania danych jako bezserwerowym potokiem Lakeflow, potok ten jest okresowo uruchamiany ponownie. Przy następnym uruchomieniu zaczyna używać nowej wersji schematu. Może minąć nawet tydzień, zanim nowe lub zmienione pola pojawią się w tabeli pozyskiwania danych.
Aby dowiedzieć się, jak potok obsługuje rekordy, które nie odpowiadają aktualnie używanemu schematowi, zobacz Dekodowanie danych za pomocą schematów.
Ładowanie i uzupełnianie danych historycznych
Parametr ingestion_config służy do konfigurowania sposobu przechwytywania i przechowywania danych strumienia na potrzeby trenowania i obsługi.
Dostęp do strumienia jest regulowany przez tabelę pozyskiwania danych:
-
SELECTw tabeli pozyskiwania danych przyznaje uprawnienia do odczytu strumienia Stream. -
MANAGEw tabeli pozyskiwania udziela dostępu do usuwania.
Aby uzyskać więcej informacji na temat uprawnień tabel, zobacz Tabela i Informacje o uprawnieniach w Unity Catalog.
Potok pozyskiwania danych
Po utworzeniu strumienia usługa Databricks uruchamia zarządzany potok pozyskiwania danych, który stale odczytuje komunikaty z tematu Kafka i zapisuje je w tabeli Delta (tabeli pozyskiwania danych). Potok danych rozpoczyna przetwarzanie od najnowszego offsetu Kafki i działa nieprzerwanie, przechwytując tylko nowe komunikaty, które napływają po utworzeniu strumienia. Ta tabela pozyskiwania służy do trenowania z funkcjami przesyłania strumieniowego. Po usunięciu strumienia jego potok przetwarzania i tabela przetwarzania również zostaną usunięte.
Miejsce docelowe importu
ingestion_destination określa trzyczęściową nazwę tabeli Delta, do której są zapisywane dane strumieniowe.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Schemat tabeli pozyskiwania danych
Tabela pozyskiwania danych zawiera dane wiadomości wraz z kolumnami metadanych:
| Column | Typ | Description |
|---|---|---|
key |
Różni się (od key_schema) |
Klucz komunikatu platformy Kafka ustrukturyzowany zgodnie z podanym schematem. |
value |
Różni się (od payload_schema) |
Wartość komunikatu Kafka (payload), ustrukturyzowana zgodnie z podanym schematem. |
stream_record_timestamp |
TIMESTAMP |
Sygnatura czasowa rekordu. Dla danych z uzupełnianiem wprzód jest to znacznik czasu przyjęcia przez brokera Kafka. W przypadku danych uzupełniających są one dostarczane przez klienta. |
kafka_topic |
STRING |
Temat Kafka, z którego odczytano rekord. |
kafka_partition |
INT |
Partycja Kafka, z której rekord został odczytany. |
kafka_offset |
LONG |
Przesunięcie platformy Kafka rekordu w ramach partycji. |
record_source |
STRING |
Albo "stream" (uzupełnianie do przodu ze strumienia Kafka na żywo), albo "backfill" (ze źródła uzupełniania wstecznego). |
Źródło wypełniania
Ponieważ potok uzupełniania zaczyna od najnowszego offsetu Kafka, nie przechwytuje wiadomości, które istniały przed utworzeniem strumienia. Aby zapewnić dostęp do danych historycznych na potrzeby szkolenia, skonfiguruj opcjonalne źródło uzupełniania danych historycznych.
Po skonfigurowaniu źródła uzupełniania danych usługa Databricks uruchamia jednorazowe zadanie MERGE INTO, które kopiuje wiersze uzupełnianych danych do tabeli ingestii za pomocą polecenia record_source="backfill". Funkcja MERGE jest uruchamiana dopiero po upewnieniu się, że źródło wypełniania wstecznego i strumień wypełniania dalej mają nakładające się znaczniki czasu (zobacz Nakładanie się między danymi wypełniania i strumienia na żywo). Jeśli warunek pokrywania się nie zostanie spełniony w ciągu 2 dni, operacja MERGE i tak zostanie uruchomiona, aby uniknąć blokowania w nieskończoność.
Tabela wypełniania musi zawierać kolumnę stream_record_timestamp typu TIMESTAMP w strefie czasowej UTC. Inne kolumny metadanych Kafka (kafka_topic, kafka_partition, kafka_offset) są przekazywane dalej, jeśli są obecne w źródle backfillu, a w przeciwnym razie są ustawiane na NULL.
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"
),
)
Nakładanie się między danymi wypełniania i transmisji strumieniowej na żywo
Przed wykonaniem operacji MERGE między uzupełnianiem danych historycznych a tabelą ładowania danych sprawdzanie nakładania się porównuje znaczniki czasu w obu tabelach:
-
Maksimum uzupełniania wstecznego: Maksymalna wartość
stream_record_timestampw źródle uzupełniania wstecznego. -
Minimum ładowania: Minimalna
stream_record_timestampliczba wierszy (record_source="stream") w tabeli ładowania.
Operacja MERGE jest wykonywana, gdy najnowszy znacznik czasu backfillu jest późniejszy od najwcześniejszego znacznika czasu tabeli ingestii o co najmniej 1 godzinę. To nakładanie się zapewnia, że w tabeli importu nie ma żadnych luk. Jeśli warunek pokrywania się nie zostanie spełniony w ciągu 2 dni, operacja MERGE i tak zostanie uruchomiona, aby uniknąć blokowania w nieskończoność.
Ponieważ potok ingestii zaczyna od najnowszego offsetu Kafka, rejestruje tylko komunikaty, które napływają po utworzeniu strumienia. Źródło danych do uzupełnienia historycznego musi zawierać dane, które obejmują zakres czasu pozyskiwania danych — a nie tylko sięgają do momentu utworzenia strumienia.
Jeśli na przykład tworzysz strumień o godzinie 15:00, potok wypełniania dalej rozpocznie odczytywanie komunikatów od 15:00. Źródło wypełniania musi zawierać dane ze znacznikami czasu do co najmniej 16:00 (1 godzina po rozpoczęciu wypełniania do przodu), aby spełnić nakładające się sprawdzanie. Oznacza to, że należy zaktualizować swoją tabelę uzupełnień po godzinie 16:00, aby upewnić się, że w tabeli ładowania nie ma luk.
Deduplication
Użyj deduplication_columns, aby określić ścieżki kolumn służące do identyfikowania zduplikowanych wierszy podczas pozyskiwania danych między danymi strumieniowymi backfill i forward-fill. Użyj notacji kropkowej dla zagnieżdżonych pól (na przykład "value.user_id").
Wybierz kolumny do usuwania duplikatów w oparciu o dane:
- Jeśli każdy rekord w strumieniu zawiera unikatowy identyfikator (na przykład
value.transaction_id), użyj tej kolumny do deduplikacji. - Jeśli źródło danych do uzupełnienia zawiera kolumny
kafka_partitionikafka_offset, użyj ich do jednoznacznego identyfikowania każdego rekordu. - Jeśli nie określono żadnych kolumn deduplikacji, domyślny klucz deduplikacji jest pełną kombinacją ,
keyvalueistream_record_timestamp. Nie jest to zalecane, ponieważ te ścisłe kryteria dopasowania mogą łatwo prowadzić do duplikatów.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Zarządzanie strumieniami
Pobierz strumień
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Wyświetlanie listy strumieni
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
Ustaw wartość include_schemas=True , aby uwzględnić pełne szczegóły schematu. Schematy mogą być duże i może to spowodować długotrwałą operację. Aby pobrać schematy indywidualnie, użyj polecenia get_stream.
Usuń strumień
Usunięcie strumienia usuwa również potok pozyskiwania danych i tabelę pozyskiwania danych.
Warning
Wszystkie modele lub funkcje odwołujące się do usuniętego strumienia nie będą już miały dostępu do danych źródłowych strumienia. Utwórz kopię tabeli pozyskiwania danych przed jej usunięciem, jeśli potrzebujesz tych danych, ale nie potrzebujesz już strumienia.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Przykładowy notatnik
Pełny przykład tworzenia obiektu Stream, definiowania funkcji strumieniowych i wdrażania w punkcie końcowym serwującym znajduje się w następującym notatniku: