Nawiązywanie połączenia z platformą Apache Kafka

Na tej stronie opisano, jak używać Apache Kafka jako źródła lub elementu docelowego podczas uruchamiania obciążeń Structured Streaming w usłudze Azure Databricks.

Aby uzyskać więcej informacji na temat platformy Kafka, zobacz dokumentację platformy Apache Kafka.

Odczytywanie danych z platformy Kafka

Użyj formatu kafka, aby skonfigurować połączenia z Kafka. Poniżej przedstawiono przykład odczytu przesyłania strumieniowego:

Python

df = (spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "latest")
  .load()
)

Scala

val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "latest")
  .load()

SQL

CREATE OR REFRESH STREAMING TABLE <table_name> AS
SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<server:ip>',
  subscribe => '<topic>'
);

Azure Databricks obsługuje również odczyty wsadowe z Kafka, jak pokazano w poniższym przykładzie:

Python

df = (spark.read
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load()
)

Scala

val df = spark.read
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("subscribe", "<topic>")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load()

SQL

SELECT * FROM read_kafka(
  bootstrapServers => '<server:ip>',
  subscribe => '<topic>',
  startingOffsets => 'earliest',
  endingOffsets => 'latest'
);

W przypadku ładowania przyrostowego wsadowego usługa Databricks zaleca używanie platformy Kafka z usługą Trigger.AvailableNow. Zobacz AvailableNow: Przyrostowe przetwarzanie wsadowe.

W środowisku Databricks Runtime 13.3 LTS i nowszym Azure Databricks udostępnia również funkcję SQL do odczytywania danych platformy Kafka. Przesyłanie strumieniowe za pomocą języka SQL jest obsługiwane tylko w potokach lakeflow lub w tabelach przesyłania strumieniowego w usłudze Databricks SQL. Zobacz read_kafka funkcji wartości tabeli.

Konfiguracja czytnika Kafka Structured Streaming

W przypadku zapytań wsadowych i przesyłanych strumieniowo należy ustawić serwery bootstrap dla źródła platformy Kafka przy użyciu następującej opcji:

Key Wartość Opis
kafka.bootstrap.servers Rozdzielona przecinkami lista host:port Serwery rozruchowe klastra Kafka

Aby ustawić tematy subskrypcji, należy określić jedną z następujących opcji:

Option Wartość Opis
subscribe Rozdzielona przecinkami lista tematów. Lista tematów do subskrybowania.
subscribePattern Java ciąg wyrażeń regularnych. Wzorzec używany do subskrybowania tematów.
assign Ciąg JSON {"topicA":[0,1],"topic":[2,4]}. Specyficzne topicPartitions do zużycia.

Zobacz Kafka , aby uzyskać pełną listę dostępnych opcji.

Schemat wierszy Kafka

Czytnik Structured Streaming dla platformy Kafka zwraca wiersze o następującym schemacie:

Kolumna Typ
key binary
value binary
topic string
partition int
offset long
timestamp timestamp
timestampType int

Obiekty key i value są zawsze deserializowane jako tablice bajtów przy użyciu ByteArrayDeserializer. Użyj operacji ramki danych (takich jak cast("string") lub from_avro), aby jawnie deserializować klucze i wartości.

Zapisywanie danych na platformie Kafka

Poniżej przedstawiono przykład przesyłania strumieniowego zapisu na platformie Kafka:

Python

(df.writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .start()
)

Scala

df.writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .start()

Azure Databricks obsługuje również semantykę zapisu wsadowego do miejsc docelowych danych Kafka, jak pokazano w poniższym przykładzie:

Python

(df.write
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .save()
)

Scala

df.write
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:ip>")
  .option("topic", "<topic>")
  .save()

Konfigurowanie składnika zapisu przesyłania strumieniowego Apache Kafka

Ważna

Środowisko Databricks Runtime 13.3 LTS i nowsze zawiera bardziej aktualną wersję biblioteki kafka-clients, która domyślnie umożliwia idempotentne zapisy. Jeśli wyjście Kafka używa wersji 2.8.0 lub starszej ze skonfigurowanymi ACL, ale bez włączenia IDEMPOTENT_WRITE, zapis kończy się niepowodzeniem z komunikatem o błędzie org.apache.kafka.common.KafkaException:Cannot execute transactional method because we are in an error state.

Rozwiąż błąd, uaktualniając do wersji 2.8.0 lub nowszej Kafki, albo ustawiając .option(“kafka.enable.idempotence”, “false”) podczas konfigurowania pisarza Structured Streaming.

Poniżej przedstawiono typowe opcje zapisu na platformie Kafka:

Key Wartość Wartość domyślna Opis
kafka.boostrap.servers Rozdzielona przecinkami lista <host:port> none Required. Konfiguracja Kafka bootstrap.servers.
topic STRING nie ustawiono Optional. Ustawia topic dla wszystkich wierszy, które mają zostać zapisane. Ta opcja zastępuje dowolną kolumnę tematu, która istnieje w danych.
includeHeaders BOOLEAN false Optional. Określenie, czy uwzględnić nagłówki Kafka w wierszu.

Zobacz Kafka sink, aby uzyskać pełną listę dostępnych opcji.

Schemat modułu zapisywania platformy Kafka

Podczas zapisywania danych na platformie Kafka podana ramka danych może zawierać następujące pola:

Nazwa kolumny Wymagane lub opcjonalne Typ
key opcjonalny STRING lub BINARY
value required STRING lub BINARY
headers opcjonalny ARRAY
topic opcjonalne (ignorowane, jeśli topic jest ustawiona jako opcja pisarza) STRING
partition opcjonalny INT

Uwierzytelnianie

Azure Databricks obsługuje wiele metod uwierzytelniania dla Kafki, w tym poświadczenia usługi Unity Catalog, SASL/SSL oraz opcje specyficzne dla chmury dla usług AWS MSK, Azure Event Hubs i Google Cloud Managed Kafka. Zobacz Uwierzytelnianie.

Pobieranie metryk platformy Kafka

Aby monitorować opóźnienie względem platformy Kafka dla zapytania strumieniowego, użyj metryk avgOffsetsBehindLatest, maxOffsetsBehindLatest i minOffsetsBehindLatest. Te metryki podają średnie, maksymalne i minimalne opóźnienie offsetu we wszystkich subskrybowanych partycjach tematów względem najnowszych offsetów w Kafka. Zobacz Interaktywne odczytywanie metryk.

Uwaga / Notatka

W środowisku Databricks Runtime 17.1 lub nowszym najnowsze przesunięcia w Kafka są pobierane po zakończeniu poszczególnych mikro-serii. W tematach, które stale odbierają dane, miary zaległości mogą wskazywać małe, trwałe wartości inne niż zero. Jest to oczekiwane zachowanie i nie wskazuje, że strumień się opóźnia.

W środowisku Databricks Runtime 17.0 i poniżej, najnowsze offsety Kafka są pobierane w czasie rozpoczęcia przetwarzania mikrosadowego. Metryki zaległości mogą zwracać 0, gdy zapytania strumieniowe stale zużywają wszystkie rekordy dostępne na początku mikropartii.

Aby oszacować pozostałe dane dla zapytania do odczytania, użyj estimatedTotalBytesBehindLatest metryki. Ta metryka szacuje łączną liczbę bajtów pozostałych we wszystkich subskrybowanych partycjach na podstawie partii przetworzonych w ciągu ostatnich 300 sekund. Możesz zmodyfikować przedział czasu używany dla tego oszacowania, ustawiając bytesEstimateWindowLength opcję .

Aby na przykład ustawić długość okna na 10 minut:

Python

df = (spark.readStream
  .format("kafka")
  .option("bytesEstimateWindowLength", "10m") # m for minutes, you can also use "600s" for 600 seconds
)

Scala

val df = spark.readStream
  .format("kafka")
  .option("bytesEstimateWindowLength", "10m") // m for minutes, you can also use "600s" for 600 seconds

Jeśli używasz strumienia w notatniku, możesz zobaczyć te wskaźniki w zakładce Nieprzetworzone dane na tablicy postępu zapytania streamingowego.

{
  "sources": [
    {
      "description": "KafkaV2[Subscribe[topic]]",
      "metrics": {
        "avgOffsetsBehindLatest": "4.0",
        "maxOffsetsBehindLatest": "4",
        "minOffsetsBehindLatest": "4",
        "estimatedTotalBytesBehindLatest": "80.0"
      }
    }
  ]
}

Aby uzyskać więcej informacji, zobacz Monitorowanie zapytań przesyłania strumieniowego ze strukturą w Azure Databricks.

Przykład dla platformy Kafka do usługi Delta Lake

Poniższy przykład przedstawia kompletny przepływ pracy na potrzeby przyrostowego zapisu strumieniowego z platformy Kafka do tabeli Delta Lake przy użyciu wyzwalacza availableNow. Tego podejścia można użyć w przypadku obciążeń pozyskiwania danych przyrostowych.

W tym przykładzie użyto stałego schematu JSON. W przypadku innych formatów, takich jak Avro lub Protobuf, użyj polecenia from_avro lub from_protobuf. Można również zintegrować z rejestrem schematów. Zobacz Przykład z rejestrem schematów.

Python

from pyspark.sql.functions import from_json, col

# Define simple JSON schemas for key and value
key_schema = "user_id STRING"
value_schema = "event_type STRING, event_ts TIMESTAMP"

# Configure Kafka options with service credentials
kafka_options = {
  "kafka.bootstrap.servers": "<bootstrap-server>:9092",
  "subscribe": "<topic-name>",
  "databricks.serviceCredential": "<service-credential-name>",
}

# Read from Kafka and parse JSON
parsed_df = (spark.readStream
  .format("kafka")
  .options(**kafka_options)
  .load()
  .select(
    from_json(col("key").cast("string"), key_schema).alias("key"),
    from_json(col("value").cast("string"), value_schema).alias("value")
  )
  .select("key.*", "value.*")
)

# Write to Delta table
query = (parsed_df.writeStream
  .format("delta")
  .option("checkpointLocation", "/path/to/checkpoint")
  .trigger(availableNow=True)
  .toTable("catalog.schema.events_table")
)

query.awaitTermination()

Scala

import org.apache.spark.sql.functions.{from_json, col}
import org.apache.spark.sql.streaming.Trigger
import org.apache.spark.sql.types.StructType

// Define JSON schemas for key and value
val keySchema = "user_id STRING"
val valueSchema = "event_type STRING, event_ts TIMESTAMP"

// Configure Kafka options with service credentials
val kafkaOptions = Map(
  "kafka.bootstrap.servers" -> "<bootstrap-server>:9092",
  "subscribe" -> "<topic-name>",
  "databricks.serviceCredential" -> "<service-credential-name>"
)

// Read from Kafka and parse JSON
val parsedDF = spark.readStream
  .format("kafka")
  .options(kafkaOptions)
  .load()
  .select(
    from_json(col("key").cast("string"), StructType.fromDDL(keySchema)).alias("key"),
    from_json(col("value").cast("string"), StructType.fromDDL(valueSchema)).alias("value")
  )
  .select("key.*", "value.*")

// Write to Delta table
val query = parsedDF.writeStream
  .format("delta")
  .option("checkpointLocation", "/path/to/checkpoint")
  .trigger(Trigger.ProcessingTime("10 seconds"))
  .toTable("catalog.schema.events_table")

query.awaitTermination()

SQL

-- Create a streaming table from Kafka using read_kafka
CREATE OR REFRESH STREAMING TABLE catalog.schema.events_table AS
SELECT
  key::string:user_id AS user_id,
  value::string:event_type AS event_type,
  to_timestamp(value::string:event_ts) AS event_ts
FROM STREAM read_kafka(
  bootstrapServers => '<bootstrap-server>:9092',
  subscribe => '<topic-name>',
  serviceCredential => '<service-credential-name>'
);

Uwaga / Notatka

W bezserwerowych zasobach obliczeniowych Databricks wyzwalacz availableNow jest zalecany w przypadku przyrostowego przesyłania strumieniowego. W przypadku ciągłego przesyłania strumieniowego o niskich opóźnieniach użyj trybu ciągłego potoków Lakeflow. Zobacz wyzwalacze funkcji Structured Streaming, aby uzyskać pełną listę obsługiwanych opcji.