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.
Często zadawane pytania dotyczące korzystania z platformy Kafka z Azure Databricks.
Dlaczego otrzymuję błąd, że opcja platformy Kafka nie jest obsługiwana lub nie jest rozpoznawana?
Ten błąd występuje, jeśli zapomnisz użyć prefiksu kafka. podczas ustawiania opcji konfiguracji klienta platformy Kafka. Wszystkie opcje przekazywane bezpośrednio do klienta platformy Kafka muszą być poprzedzone prefiksem kafka.:
Poniższy kod przedstawia nieprawidłowe opcje, które nie mają prefiksu kafka. :
.option("security.protocol", "SASL_SSL")
.option("sasl.mechanism", "PLAIN")
Poniższy kod przedstawia poprawne opcje:
.option("kafka.security.protocol", "SASL_SSL")
.option("kafka.sasl.mechanism", "PLAIN")
Opcje łącznika platformy Spark Kafka (na przykład subscribe, startingOffsets, maxOffsetsPerTrigger) nie wymagają prefiksu. Aby uzyskać pełną listę opcji, zobacz Kafka.
Dlaczego otrzymuję błąd dotyczący zacienionych klas platformy Kafka?
Azure Databricks wymaga użycia zacienionych klas platformy Kafka (poprzedzonych kafkashaded. lub shadedmskiam.). Jeśli widzisz błędy, takie jak RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED, należy użyć zacienionych nazw klas:
-
org.apache.kafka.*klasy wymagają prefiksukafkashaded.. Przykład:kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule -
software.amazon.msk.*klasy wymagają prefiksushadedmskiam.. Przykład:shadedmskiam.software.amazon.msk.auth.iam.IAMLoginModule
Dlaczego otrzymuję błąd TimeoutException podczas nawiązywania połączenia z Kafka?
Typowe przyczyny:
- Łączność sieciowa: klaster obliczeniowy nie może nawiązać połączenia z brokerami platformy Kafka. Sprawdź reguły zapory, grupy zabezpieczeń i konfiguracje VPC.
-
Nieprawidłowe serwery bootstrap: sprawdź, czy
kafka.bootstrap.serversnazwa hosta i port są poprawne. - Rozwiązywanie nazw DNS: Sprawdź, czy nazwy hostów brokera Kafka mogą być rozwiązywane w sieci Azure Databricks.
- Problemy z protokołem SSL/TLS: w przypadku korzystania z protokołu SSL sprawdź, czy certyfikaty są poprawnie skonfigurowane.
W przypadku konfiguracji Private Link lub peeringu VPC sprawdź, czy skonfigurowano prawidłowe trasy sieciowe.
Czy powinienem używać trybu partii czy strumieniowego dla systemu Kafka?
Zależy to od przypadku użycia:
- Tryb przesyłania strumieniowego (): Używaj, gdy potrzebujesz ciągłego przetwarzania danych lub pozyskiwania danych przy niskich opóźnieniach.
-
Tryb wsadowy (
spark.read): Używany do jednorazowego ładowania danych, ich uzupełniania lub debugowania. Wymaga zarównostartingOffsets, jak iendingOffsets.
Zobacz Konfigurowanie interwałów wyzwalacza strukturalnego przesyłania strumieniowego, aby uzyskać szczegółowe informacje na temat konfigurowania interwałów wyzwalaczy, takich jak AvailableNow, ProcessingTimei tryb czasu rzeczywistego.
Czy mogę odczytać z wielu tematów platformy Kafka w jednym strumieniu?
Tak, możesz użyć:
-
subscribe: Podaj rozdzielaną przecinkami listę tematów, na przykład.option("subscribe", "topic1,topic2"). -
subscribePattern: użyj wzorca wyrażenia regularnego Java, aby dopasować nazwy tematów, na przykład.option("subscribePattern", "topic-.*").
Jak używać platformy Kafka z potokami Lakeflow?
Potoki Lakeflow mają wbudowaną obsługę źródeł Kafka.
Możesz zdefiniować tabelę strumieniową, która odczytuje dane z platformy Kafka, jak w poniższym kodzie:
Python
import dlt
@dlt.table
def kafka_bronze():
return (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.load()
)
SQL
CREATE OR REFRESH STREAMING TABLE kafka_bronze AS
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:port>',
subscribe => '<topic>'
);
Zobacz Ładowanie danych w potokach, aby uzyskać więcej informacji o źródłach strumieniowych w potokach Lakeflow.
Jak zdeserializować kolumny klucza i wartości w Kafka?
Kolumny key i value są zwracane jako typ BINARY. Użyj operacji DataFrame do deserializacji ich na podstawie formatu danych.
-
Dane ciągu: służy
cast("string")do konwertowania danych binarnych na ciąg. -
Dane JSON: użyj polecenia
from_json()po rzutowaniu na ciąg. Zobaczfrom_jsonfunkcję. -
Dane Avro: Użyj
from_avro()do deserializacji danych zakodowanych w formacie Avro. Zobacz Odczytywanie i zapisywanie przesyłanych strumieniowo danych Avro. -
Bufory protokołu: służy
from_protobuf()do deserializacji danych protobuf. Zobacz Bufory protokołu odczytu i zapisu.
Dlaczego otrzymuję błąd zapisu idempotentnego?
Ś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 klaster Kafka używa wersji 2.8.0 lub niższej ze skonfigurowanymi listami ACL, ale bez włączonego IDEMPOTENT_WRITE, zapis kończy się niepowodzeniem: 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.
Co to jest KAFKA_DATA_LOSS_ERROR i jak mogę rozwiązać ten problem?
Ten błąd występuje, gdy źródło platformy Kafka wykryje, że przesunięcia przechowywane w punkcie kontrolnym nie są już dostępne na platformie Kafka, zwykle dlatego, że:
- Strumień został wstrzymany dłużej niż okres retencji w Kafka.
- Dane tematu Kafka zostały usunięte lub temat został ponownie utworzony.
- Broker platformy Kafka doświadczył utraty danych.
Aby rozwiązać:
-
Jeśli utrata danych jest akceptowalna: ustaw opcję
.option("failOnDataLoss", "false")zezwalania strumieniowi na kontynuowanie od najwcześniejszego dostępnego przesunięcia. -
Jeśli utrata danych nie jest akceptowalna: zresetuj punkt kontrolny i ponownie przeprocesuj z
earliestprzesunięć lub przywróć brakujące dane platformy Kafka.
Zobacz warunek błędu KAFKA_DATA_LOSS aby uzyskać więcej informacji.
Jak kontrolować szybkość odczytywania danych z platformy Kafka?
Użyj opcji maxOffsetsPerTrigger, aby ograniczyć liczbę przesunięć (w przybliżeniu liczbę rekordów) przetwarzanych na mikropartie. Pomaga to zapobiegać dużym partiom, które mogą przeciążyć przetwarzanie dalsze lub powodować problemy z pamięcią podczas nadrabiania zaległości w przetwarzaniu.
Python
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.option("maxOffsetsPerTrigger", 10000)
.load()
)
Scala
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:port>")
.option("subscribe", "<topic>")
.option("maxOffsetsPerTrigger", 10000)
.load()
SQL
SELECT * FROM STREAM read_kafka(
bootstrapServers => '<server:port>',
subscribe => '<topic>',
maxOffsetsPerTrigger => '10000'
);
Alternatywnie użyj opcji, takich jak minPartitions lub maxRecordsPerPartition , aby kontrolować liczbę partycji platformy Spark tworzonych dla każdej partii.
Jak mogę monitorować, jak daleko mój strumień jest opóźniony względem najnowszych offsetów w Kafka?
Użyj metryk avgOffsetsBehindLatest, maxOffsetsBehindLatest, i minOffsetsBehindLatest dostępnych w postępie zapytania przesyłania strumieniowego. Ten raport przedstawia liczbę przesunięć za najnowszym dostępnym przesunięciem strumienia we wszystkich subskrybowanych partycjach tematu. Zobacz Monitorowanie zapytań przesyłania strumieniowego ze strukturą w Azure Databricks.
Można również użyć estimatedTotalBytesBehindLatest do oszacowania łącznej liczby bajtów danych, które nie zostały jeszcze przetworzone.
Dlaczego metryki opóźnienia przesunięcia platformy Kafka pokazują trwałe wartości inne niż zero po uaktualnieniu do środowiska Databricks Runtime 17.1?
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.
Jeśli wartości są duże lub stale rosną, strumień może nie być na bieżąco z danymi przychodzącymi. Zobacz Monitorowanie zapytań przesyłania strumieniowego ze strukturą w Azure Databricks.
Dlaczego inicjowanie strumienia Kafki jest powolne?
Strumienie Kafka wymagają czasu:
- Połącz się z klastrem platformy Kafka i pobierz metadane.
- Odkryj partycje tematów.
- Pobieranie początkowych offsetów.
W przypadku klastrów lokalnych lub zdalnych platformy Kafka opóźnienie sieci może znacząco wpłynąć na czas inicjowania. Jeśli uruchamiasz wyzwalane/zaplanowane potoki z częstymi ponownym uruchamianiem, rozważ użycie trybu ciągłego przesyłania strumieniowego, aby uniknąć wielokrotnego inicjowania obciążenia.
Dlaczego dodanie większej liczby egzekutorów Spark nie zwiększa przepustowości Kafki?
Gdy brokerzy platformy Kafka staną się nasyceni, dodanie kolejnych executorów Spark zwiększa koszt bez zwiększania przepływności.
Znaki, że Kafka jest wąskim gardłem:
- Przepustowość stabilizuje się pomimo dodania dodatkowych rdzeni.
- Wysokie zużycie CPU lub sieci brokera Kafka.
- Zadania platformy Spark są wykonywane szybko, ale czekają na nowe dane.
Aby rozwiązać ten problem, przeprowadź skalowanie klastra platformy Kafka przez dodanie brokerów lub zwiększenie liczby partycji w celu dystrybucji obciążenia.
Jak zoptymalizować koszty i wykorzystanie zasobów obliczeniowych na potrzeby przesyłania strumieniowego platformy Kafka?
W przypadku trybów mikrosadowych i AvailableNow:
- Odpowiedni rozmiar klastra: Monitoruj metryki i ustaw odpowiedni stały rozmiar klastra pod kątem szczytowego obciążenia.
-
Użyj polecenia
maxOffsetsPerTrigger: Ogranicz rozmiary partii, aby kontrolować użycie zasobów podczas skoków obciążenia. - Unikaj skalowania automatycznego: zadania przesyłania strumieniowego są uruchamiane w sposób ciągły, a dodawanie lub usuwanie węzłów powoduje ponowne równoważenie obciążenia zadania.
-
Zmniejsz niesymetryczność danych: niesymetryczne partycje powodują, że niektóre zadania przetwarzają znacznie więcej danych niż inne, co prowadzi do spowolnienia ogólnego ukończenia partii i marnowania zasobów obliczeniowych na zadania bezczynne. Użyj opcji dzielenia
minPartitionsdużych partycji platformy Kafka na mniejsze partycje platformy Spark w celu bardziej zrównoważonego przetwarzania.
W przypadku trybu czasu rzeczywistego ustalanie rozmiaru zasobów obliczeniowych jest szczególnie ważne, ponieważ zadania mogą pozostawać bezczynne podczas oczekiwania na dane. Kluczowe kwestie:
- Ustaw
maxPartitionstak, aby każde zadanie obsługiwało wiele partycji platformy Kafka, aby zmniejszyć obciążenie. - Dostosuj
spark.sql.shuffle.partitionsdo zadań z intensywnym przetwarzaniem shuffle.
Zobacz Ustalanie rozmiaru zasobów obliczeniowych , aby uzyskać wskazówki dotyczące określania rozmiaru klastrów w trybie czasu rzeczywistego.
Dlaczego strumień nie zwraca żadnych rekordów, mimo że dane istnieją w topicu?
Typowe przyczyny:
-
Nieprawidłowe
startingOffsetsustawienie: wartość domyślna tolatest, która odczytuje tylko nowe dane odbierane po uruchomieniu strumienia. UstawstartingOffsetsnaearliest, aby odczytać istniejące dane. - Nieprawidłowa nazwa tematu: Sprawdź, czy subskrybujesz właściwy temat.
- Problemy z uwierzytelnianiem: Strumień mógł się połączyć pomyślnie, ale brak mu uprawnień do odczytu z kanału. Sprawdź listy ACL platformy Kafka.
-
Wygaśnięcie przesunięcia: Jeśli strumień został zatrzymany przez dłuższy okres, a przesunięcia w punkcie kontrolnym wygasły (zostały usunięte przez retencję Kafka), konieczne może być zresetowanie punktu kontrolnego lub dostosowanie
failOnDataLoss.