Akış ayarlama

Ö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 bkz. Alma ve geri doldurma .

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-client Python paketi sürüm 0.16.0 veya üzeri yüklü olmalıdı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ı: Akış platformunu (örneğin, Kafka) ve kaynağa özgü ayrıntıları (Kafka için konu aboneliği gibi) 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. Ayrıntılar için bkz. Alma ve geri doldurma .
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"
        ),
    ),
)

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. Bkz. Sunucusuz işlem üzerinde akış ve Apache Kafka’ya bağlanma.

AWS tarafından yönetilen akış (Amazon MSK) için bkz. Amazon MSK'ye sunucusuz özel bağlantı. Kafka kimlik doğrulama seçenekleri hakkında ayrıntılı bilgi için bkz. Kimlik doğrulaması.

Authentication

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.

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

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

SASL

SASL kimlik doğrulaması (hem SASL/SCRAM hem de SASL/PLAIN) önizleme sırasında desteklenmez.

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

Şema yapılandırması

JSON Şeması biçimini kullanarak ileti anahtarlarının ve değerlerinin yapısını tanımlayın. 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.

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

Bir anahtar veya yük için şema sağlanmazsa, basit bir dize olarak değerlendirilir.

Alma 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:

  • SELECT veri alımı tablosunda Stream'e okuma erişimi verir.
  • MANAGE alma 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, Kafka konusundan gelen iletileri sürekli olarak okuyan ve bunları bir Delta tablosuna (veri alım tablosu) yazan yönetilen bir veri alım işlem hattı başlatır. İşlem hattı en son Kafka offset’inden başlar ve kesintisiz olarak çalışarak yalnızca akış oluşturulduktan sonra gelen yeni iletileri 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ı

Alma tablosu, meta veri sütunlarıyla birlikte ileti verilerini içerir:

Column Türü Description
key Değişir (key_schema öğesinden itibaren) Kafka ileti anahtarı, sağladığınız şemaya göre yapılandırılmıştır.
value Değişir (payload_schema öğesinden itibaren) Kafka ileti değeri (yük), sağladığınız şemaya göre yapılandırılmıştır.
stream_record_timestamp TIMESTAMP Kayıt zaman damgası. İleriye doğru doldurma verileri için bu, Kafka aracısı alma zaman damgasıdır. Geri doldurma verileri için bu müşteri tarafından sağlanır.
kafka_topic STRING Kaydın tüketildiği Kafka konusu.
kafka_partition INT Kaydın tüketildiği Kafka bölümü.
kafka_offset LONG Bölümü içindeki kaydın Kafka uzaklığı.
record_source STRING Ya "stream" (canlı Kafka akışından ileri doldurma) ya da "backfill" (geri doldurma kaynağından).

Geri doldurma kaynağı

İleriye doğru doldurma işlem hattı en son Kafka uzaklığından başladığından, akış oluşturulmadan önce var olan iletileri 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 Kafka meta veri sütunları (kafka_topic, kafka_partition, kafka_offset), geri doldurma kaynağında mevcutsa olduğu gibi aktarılır; aksi takdirde NULL olarak ayarlanır.

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_timestamp değer.
  • Minimum alma: Alma tablosundaki satırların minimum stream_record_timestamp sayı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 alım hattı en son Kafka offset’inden başladığı için yalnızca akış oluşturulduktan sonra gelen iletileri 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_partition ve kafka_offset sü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, value ve stream_record_timestamp sü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önetme

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:

Akış Özellik Görünümleri Hızlı Başlangıç Not Defteri

Dizüstü bilgisayar al