Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Önemli
Bu özellik Genel Önizleme aşamasındadır. Çalışma alanı yöneticileri Bu özelliğe erişimi Önizlemeler sayfasından denetleyebilir. Bkz. Azure Databricks önizlemelerini yönetme.
Akış, Apache Kafka gibi bir dış akış veri kaynağını temsil eder. Akışlar bağlantı ayrıntılarını, kimlik doğrulamasını, şemaları ve alım yapılandırmasını depolar. Bir akış oluşturulduktan sonra, gerçek zamanlı akış özellikleri oluşturmak için Özellik Görünümü tanımlarını kullanarak bu akışa başvurabilirsiniz.
Akışlar üç bölümlü adlara (catalog.schema.stream_name) sahiptir. Akışa erişim, ilişkili alım tablosu tarafından yönetilir. Ayrıntılar için Veri alımı ve geriye dönük doldurma bölümüne bakın.
Requirements
- Not defteri komutlarını çalıştırmak için: sunucusuz veya Databricks Runtime 17.0 ML veya üzerini çalıştıran klasik bir işlem kümesi.
-
feature-engineering-clientPython paket sürüm 0.17.0 veya üzeri yüklenmelidir.
Akış kaynaklarına bağlanma
Akış özelliklerini tanımlamadan önce, Kafka broker’ınıza bir Lakeflow akış işlem hattı bağlantısı oluşturun ve bunu test edin. Feature Store, sunucusuz SDP'ye dayanır, bu da klasik hesaplamanızı (broker veya uç nokta) Databrick'in sunucusuz hesaplamasına bağlamak için bir mekanizmaya ihtiyacınız olduğu anlamına gelir. Bu, privatelink gibi ürünlerle veya klasik hesaplamanızın halka açık internetten erişilebilir olmasına izin vererek yapılır.
Akış oluşturma
Yeni bir Stream oluşturmak için kullanın create_stream() . Akış dört yapılandırma bileşeni gerektirir:
- Kaynak yapılandırması: Yayın platformunu ve kaynağa özgü detayları, örneğin Kafka kaynağı için konu aboneliğini belirtir.
- Bağlantı yapılandırması: Önyükleme sunucuları ve kimlik bilgileri dahil olmak üzere akış platformuna bağlanmayı ve kimlik doğrulamayı belirtir.
- Şema yapılandırması: İleti anahtarlarının ve değerlerinin yapısını tanımlar.
- İçe aktarma yapılandırması: Akış verilerinin nereden ve nasıl içe aktarılacağını belirtir. Detaylar için Yutma ve Doldurma bölümlerine bakınız.
Kaynağa özgü source_config, bağlantı kurulumu ve eksiksiz bir create_stream() örneği için Apache Kafka bölümüne bakın.
Şema ve veri alma seçenekleri kaynaklar arasında ortaktır.
Apache Kafka
Apache Kafka'dan akış almak için, kaynak yapılandırması olarak KafkaStreamConfig ve kimlik doğrulama için bir Unity Catalog bağlantısı kullanın. Kafka bağlantısı için Sunucusuz işlemde akış ve Apache Kafka’ya bağlanma konularına bakın.
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 abonelik modları
Abonelik modu, Stream'in kullanılacak Kafka konularını nasıl seçtiğini belirtir. Üç mod desteklenir:
| Mode | Description | Example |
|---|---|---|
subscribe |
Konu adlarının virgülle ayrılmış listesi | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Java regex deseni eşleşen konu adları | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
Konu bölümü atamalarını belirten JSON | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Kafka doğrulaması
Unity Kataloğu bağlantısı (önerilir)
Kafka kümenizde kimlik doğrulaması yapmak için Unity Kataloğu bağlantısı kullanın. Bu, yönetilen kimlik doğrulaması için önerilen yaklaşımdır. Bağlantı oluşturmak için bkz. Bağlantı oluşturma. Akışın oluşturucusunda bağlantıda USE CONNECTION olmalıdır. Stream’i kaynak olarak kullanan herhangi bir kullanıcının özellik oluşturabilmesi için, bağlantı üzerinde USE CONNECTION da bulunmalıdır.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
Bağlantı hem IAM (hizmet yetkinliği) hem de SASL kimlik doğrulamasını destekler.
IAM (hizmet sertifikası)
Örneğin, Amazon MSK'ya IAM ile bağlanmak için Unity Kataloğu hizmet kimlik bilgileriyle kimlik doğrulama yapın. Hizmet kimlik bilgileri oluşturmak için bkz. Hizmet kimlik bilgileri oluşturma. Hizmet kimlik adı adını şu credential seçenekle ayarlayın:
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>'
)
Bağlantıya USE CONNECTION ek olarak, hizmet kimliğini kullanan kimlikler de bu bağlantıda olmalıdır ACCESS . Referans verilen hizmet kimlik bilgisi üzerinde ACCESS iznini Stream’in oluşturucusuna ve Stream ile özellikleri somutlaştıran herhangi bir kimliğe verin. Bkz Hizmet kimlik bilgilerini kullanarak bir dış bulut hizmetine erişim izni verme.
SASL
SASL doğrulaması kullanıcı adı ve şifre kullanır. Aşağıdakilerden birine ayarlayın sasl_mechanism :
PLAINSCRAM-SHA-256SCRAM-SHA-512
Kimlik bilgilerini user ve password seçenekleriyle belirtin. Bağlantı bu kimlik bilgilerini güvenli bir şekilde saklar.
Aşağıdaki örnekte SASL/SCRAM kullanılmaktadır. SASL/PLAIN için, sasl_mechanism değerini PLAIN olarak ayarlayın.
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
sasl_mechanism 'SCRAM-SHA-512',
user '<username>',
password '<password>'
)
Doğrudan mTLS
Doğrudan mTLS kimlik doğrulaması için, Unity Catalog biriminde depolanan keystore ve truststore dosyalarını, parolaları Databricks secret scope’ları üzerinden başvurulacak şekilde sağlayın. Kafka ile SSL kimlik doğrulaması hakkında daha fazla bilgi için bkz. Azure Databricks Kafka'ya bağlanmak için SSL kullanma.
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"
),
),
)
Şema konfigürasyonu
Mesaj anahtarlarının ve değerlerin yapısını tanımlayın ki alım ve özellik tanımları bireysel alanları okuyabilsin. Kafka kaynakları için, payload_schema Kafka ileti değerine ( value Kafka'nın anahtar-değer modelindeki) ve key_schema Kafka ileti anahtarına karşılık gelir.
payload_schema veya key_schema değerlerinden en az biri belirtilmelidir.
Her biri, SchemaConfig kaynağın mesajlarını serileştirme şekline uyan üç formattan birini kabul eder: json_schema, avro_schema, veya proto_schema. Bir anahtar veya yük için şema sağlanmazsa, basit bir dize olarak değerlendirilir.
Bu bölümdeki kod örnekleri, DirectSchemas ile satır içinde tanımlanan şemaları kullanır; burada şema bir dize olarak sağlanır. Şemaları harici şema kayıt defteri kullanarak yönetmek için detaylar için Schema registrine bakınız.
JSON şeması
json_schema için bir JSON Şeması dizesi sağlayın.
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 şeması
avro_schema için bir Avro şeması dizesi sağlayın. Avro mantıksal türleri desteklenir; bunlara timestamp-millis, date ve decimal dahildir.
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 şeması
.proto öğesine, ProtoSchemaSpecproto_schema kaynak metni ve yük mesajı adıyla birlikte sağlayın.
ProtoSchemaSpec öğesini databricks.feature_engineering.entities öğesinden içeri aktar.
message_name, com.example.Event metninde bildirilen package dahil olmak üzere tam nitelikli ileti adı olmalıdır (örneğin, Event, .proto değil). Hem proto2 hem de proto3 sözdizimi desteklenir.
google.protobuf.Timestamp ve skaler wrapper türleri (StringValue, Int32Value, vb.) desteklenir ve içe aktarmaları otomatik olarak çözülür. Diğer bilinen türler, örneğin Duration, Struct, ve Any, reddedilir; bu değerler desteklenen bir skaler veya mesaj olarak kodlanır.
fixed32 ve fixed64 skaler türleri ile dize olmayan anahtarlar içeren map de desteklenmez.
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",
)
),
)
Verileri şemalarla çözme
Databricks, her mesajı Spark'ın from_json, from_avro, ve from_protobuf fonksiyonlarıyla çözer. Aşağıdaki davranışlar, şemayı satırda ilan etseniz de ya da şema kayıt defterinden çözerseniz de geçerlidir:
- Bozuk kayıtlar. Çözümleme,
PERMISSIVEmodunu kullanır; bu nedenle şemasıyla eşleşmeyen bir kayıt, akışın başarısız olmasına neden olmak yerine null değer olarak çözümlenir. - Avro sendikaları. Birden fazla kayıt türünün birleşimi, her kayıt türü için bir alan olan ve her biri Avro kaydının adını taşıyan bir yapıya çözümlenir.
- Protobuf tipleri. İşaretsiz tam sayılar daha geniş bir işaretli türe (örneğin,
uint32içinBIGINTveuint64içinDECIMAL(20,0)), enum alanları adlarının dize gösterimine ve skaler sarmalayıcı türler (örneğin,StringValueveInt32Value) sarmalanan türün null değer alabilen bir sütununa çözülür.
Şema kayıt defteri
Şema kayıt defterleri, akış üreticileri ve tüketicilerinin kullandığı şemaları saklar ve sürümlendirir; bu şemalar geliştikçe uyumluluk kurallarını uygular. Harici bir şema kaydı yapılandırıldığında, Feature Store kayıt defterinden şemayı okur ve akış mesajını çözmek için kullanır. Şema kayıt defteri kullanırken şemayı Akış üzerinde inline olarak tanımlamazsınız.
Şema kayıt desteği aşağıdaki sınırlamalara sahiptir:
- Sadece Kafka yayınları için destekleniyor.
- Yalnızca Confluent Schema Registry desteklenmektedir
- Sadece Avro ve Protobuf formatları desteklenmektedir. JSON mesajlarını okumak için, bunun yerine şemayı satır içinde tanımlayın. JSON şemasına bakınız.
- Her Akış, mesaj değeri için tam olarak bir Confluent konusuna, mesaj anahtarı için ise (sağlanırsa) bir Confluent konusuna bağlıdır. Birden fazla şema kaydı içeren akış konuları desteklenen bir yapılandırma değildir. Akışınız birden fazla şema içeren konulara bağlanıyorsa, belirtilen konu için şema ile eşleşmeyen kayıtlar null olarak çözülür.
Bir şema kayıt defterine bağlanın
Kafka Unity Kataloğu bağlantısında kayıt bağlantı detaylarını seçenek olarak sunun ve kayıt API gizliliğini Databricks gizli kapsamında saklayın. Akışın run-as kimliğinin gizli kapsam üzerinde READ iznine sahip olması gerekir, çünkü veri alma işlem hattı gizli diziyi çalışma zamanında okur. Bağlantı oluşturma ve yapılandırma için bkz. Bağlantı oluştur.
Kimlik doğrulama için kullanılan bağlantıya , schema_registry_api_key, ve schema_registry_api_secret seçeneklerini ekleyinschema_registry_url. Aşağıdaki örnek, Unity Kataloğu hizmet kimlik bilgisi ile aracı kuruma ve API anahtarıyla kayıt defterine kimlik doğrulama sağlayan bir Kafka bağlantısı oluşturur:
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>')
)
Kafka bağlantısındaki schema_registry_api_secret seçeneğini ve Akış üzerindeki gizli kapsamı başvurusunu aynı gizli değere ayarlayın.
Şema kayıt defteri kullanan bir akış oluşturun
schema_config olarak bir SchemaRegistryConfig geçirin. Registry API gizli anahtarına api_secret_ref ile başvurun ve mesaj değeri için özneyi ve biçimi key_schema_locator ile, mesaj anahtarı içinse payload_schema_locator ile belirtin. En az bir konum belirleyici sağlanmalıdır.
Burada Şema yapılandırma bölümündeki doğrudan şema örneklerine kıyasla farklara dikkat edin. Şema kayıt defteri kullanırken, şemayı Akış üzerinde satır içi olarak schema_config sağlanmaz. Bunun yerine, kayıttaki şemayı tanımlayan bir SchemaRegistryConfig öğesini belirtirsiniz.
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"
),
),
)
Confluent öznesi, bir şemanın sürüm geçmişinin kaydedildiği ve uyumluluğun zorunlu kılındığı adlandırılmış kapsamdır. İlgili subject kapsamın adını belirleyin, bu genellikle konu adı stratejisinden belirlenir:
-
TopicNameStrategy (varsayılan, konuyu konu adından türetir):
<topic>-valuedeğer için ve<topic>-keyanahtar için. Örneğin, konutransactionsiçin değer şeması konuyutransactions-valuekullanır. -
RecordNameStrategy (konuyu şemanın kayıt adından türetir, konudan bağımsız): tam nitelikli kayıt adı, örneğin
com.example.Payment. Bu, Avro için kaydın ad alanı ve adı ya da Protobuf için mesajın paketi ve adıdır. -
TopicRecordNameStrategy (konu ve kayıt isimlerini birleştirir):
<topic>-<fully-qualified-record-name>, örneğintransactions-com.example.Payment.
format gereklidir. Bunu, konunun nasıl serileştirildiğine uyacak şekilde SchemaLocatorFormat.FORMAT_AVRO veya SchemaLocatorFormat.FORMAT_PROTOBUF olarak ayarlayın.
Şema evrimi
Alım boru hattı, başladığında konunun mevcut şemasını çözer. Schema Registry’de subject için geriye dönük uyumlu yeni bir şema sürümü kaydettirdiğinizde, çalışan iş hattı kullanmaya başladığı sürümü kullanmayı sürdürür.
Databricks, veri alım hattını Lakeflow’un sunucusuz bir işlem hattı olarak yönettiği için, işlem hattı düzenli aralıklarla yeniden başlatılır. Bir sonraki yeniden başlatmada yeni şema versiyonunu alıyor. Yeni veya değiştirilmiş alanların tüketim tablosunda görünmesi bir haftaya kadar sürebilir.
Boru hattının şu anda kullandığı şema ile eşleşmeyen kayıtları nasıl işlediği için bkz. Şemalarla veri dekodasyonu.
Yutma ve geri doldurma
parametresi, ingestion_config akış verilerinin nasıl yakalanıp eğitim ve sunum için depolandığını yapılandırır.
Akışa erişim, alma tablosu tarafından yönetilir:
-
SELECTveri alımı tablosunda Stream'e okuma erişimi verir. -
MANAGEalma tablosunda silme erişimi verir.
Tablo ayrıcalıkları hakkında daha fazla bilgi için Tablo ve Unity Catalog ayrıcalıkları başvuru bilgileri bölümlerine bakın.
Alım işlem hattı
Bir akış oluşturulduğunda Databricks, kaynak akıştaki iletileri sürekli olarak okuyup bunları bir Delta tablosuna (veri alım tablosu) yazan yönetilen bir veri alımı işlem hattı başlatır. Boru hattı, kaynaktaki en son konumdan başlar ve sürekli çalışır, yalnızca akış oluşturulduktan sonra gelen yeni mesajları yakalar. Bu veri alım tablosu, akış özellikleriyle eğitim için kullanılır. Bir akış silindiğinde, içeri aktarma işlem hattı ve içeri aktarma tablosu da silinir.
Alma hedefi
, ingestion_destination akış verilerinin yazıldığı üç bölümlü Delta tablo adını belirtir.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Alma tablosu şeması
Alım tablosu, mesaj verilerini ve meta veri sütunlarını içerir. Ortak sütunlar her kaynak için mevcuttur; sütunlar kafka_* sadece bir Kafka akımı için mevcuttur.
| Column | Türü | Source | Description |
|---|---|---|---|
key |
Değişir (key_schema öğesinden itibaren) |
Common | Mesaj tuşu, sağladığınız şemaya göre yapılandırılmış. |
value |
Değişir (payload_schema öğesinden itibaren) |
Common | Mesaj değeri (yük), sağladığınız şemaya göre yapılandırılmıştır. |
stream_record_timestamp |
TIMESTAMP |
Common | Kayıt zaman damgası. İleriye doğru doldurma verileri için, kaynak alım zaman damgasıdır. Geri doldurma verileri için bu müşteri tarafından sağlanır. |
record_source |
STRING |
Common | Ya "stream" (canlı yayından ileri doldurma) ya "backfill" da (geri doldurma kaynağından). |
kafka_topic |
STRING |
Kafka | Kaydın tüketildiği Kafka konusu. |
kafka_partition |
INT |
Kafka | Kaydın tüketildiği Kafka bölümü. |
kafka_offset |
LONG |
Kafka | Bölümü içindeki kaydın Kafka uzaklığı. |
Geri doldurma kaynağı
İleriye doğru doldurma boru hattı kaynağın en son konumundan başladığı için, akış oluşturulmadan önce var olan mesajları yakalamaz. Eğitim için geçmiş veri kapsamı sağlamak için isteğe bağlı bir doldurma kaynağı yapılandırın.
Bir geri doldurma kaynağı yapılandırıldığında Databricks, geri doldurma satırlarını MERGE INTO ile alım tablosuna kopyalayan tek seferlik bir record_source="backfill" iş çalıştırır. MERGE yalnızca örtüşme denetleyicisi geri doldurma kaynağı ile ileri doldurma akışının çakışan zaman damgalarına sahip olduğunu onayladıktan sonra çalışır (bkz. Geri doldurma ve canlı akış verileri arasında çakışma). Örtüşme koşulu 2 gün içinde karşılanmazsa, süresiz engellemeyi önlemek için MERGE yine de çalıştırılır.
Geri doldurma tablosu, UTC saat diliminde stream_record_timestamp türünde bir TIMESTAMP sütun içermelidir. Diğer meta veri sütunları, geri doldurma kaynağında mevcutsa olduğu gibi aktarılır; aksi takdirde NULL olarak ayarlanır. Kafka için bunlar kafka_topic, kafka_partition, ve 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"
),
)
Geri doldurma ve canlı akış verileri arasında örtüşme
Geri doldurma ve alma tablosu arasında MERGE çalıştırmadan önce, çakışma denetimi iki tablodaki zaman damgalarını karşılaştırır:
-
En fazla geri doldurma: Doldurma kaynağındaki maksimum
stream_record_timestampdeğer. -
Minimum alma: Alma tablosundaki satırların minimum
stream_record_timestampsayısı (record_source="stream").
MERGE işlemi, geri doldurmanın en son zaman damgası içeri aktarma tablosunun en erken zaman damgasından en az 1 saat daha ileride olduğunda gerçekleştirilir. Bu çakışma, alma tablosunda boşluk olmamasını sağlar. Örtüşme koşulu 2 gün içinde karşılanmazsa, süresiz engellemeyi önlemek için MERGE yine de çalıştırılır.
Veri alma işlem hattı kaynaktaki en güncel konumdan başladığı için yalnızca akış oluşturulduktan sonra gelen mesajları yakalar. Geri doldurma kaynağınız, yalnızca akış oluşturma zamanına kadar değil, veri alımı zaman aralığına uzanan veriler içermelidir.
Örneğin, 15:00'te bir akış oluşturursanız, ileri doldurma işlem hattı iletileri 15:00'ten itibaren okumaya başlar. Doldurma kaynağınız, çakışma denetimini karşılamak için en az 16:00'dan (ileri doldurma başlangıcının 1 saat sonrasına) kadar zaman damgalarına sahip veriler içermelidir. Bu, alma tablosunda boşluk olmadığından emin olmak için geri doldurma tablonuzu saat 16:00'dan sonra güncelleştirmeniz gerektiği anlamına gelir.
Deduplication
deduplication_columns öğesini, geri doldurma ve ileri doldurma akış verileri arasında içe alım sırasında yinelenen satırları tanımlamak için sütun yollarını belirtmede kullanın. İç içe alanlar için nokta gösterimini kullanın (örneğin, "value.user_id").
Verilerinize göre yinelenenleri kaldırma sütunlarını seçin:
- Akışınızdaki her kayıt benzersiz bir tanımlayıcı içeriyorsa (örneğin,
value.transaction_id), yinelenenleri kaldırma için bu sütunu kullanın. - Geri doldurma kaynağınız
kafka_partitionvekafka_offsetsütunlarını içeriyorsa, her kaydı benzersiz şekilde tanımlamak için bunları kullanın. - Yinelenenleri giderme sütunları belirtilmezse, varsayılan yinelenenleri giderme anahtarı
key,valuevestream_record_timestampsütunlarının tam kombinasyonudur. Bu katı ölçüt eşleştirmesi kolayca yinelemelere neden olabileceği için bu önerilmez.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Akışları yönet
Akış al
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Akış listele
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
Tam şema ayrıntılarını içerecek şekilde ayarlayın include_schemas=True . Şemalar büyük olabilir ve bu durum uzun süre çalışan bir işleme neden olabilir. Bunun yerine şemaları tek tek almak için kullanın get_stream.
Akışı silme
Bir akışın silinmesi, veri alımı işlem hattını ve veri alımı tablosunu da siler.
Warning
Silinen akışa başvuran model veya özellikler artık temel alınan akış verilerine erişemez. Bu verilere ihtiyacınız varsa ancak artık akışa ihtiyacınız yoksa, silme işleminden önce alma tablosunun bir kopyasını oluşturun.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Örnek defter
Stream oluşturan, akış özelliklerini tanımlayan ve bir sunum uç noktasına dağıtan uçtan uca bir örnek için aşağıdaki not defterine bakın: