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.
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.