Я пытаюсь отобразить данные Avro из Kafka на консоли с помощью DataStream и поймать исключение.
Причина: org.apache.spark.SparkException: при анализе записей обнаружены некорректные записи. Текущий режим анализа: FAILFAST. Чтобы обрабатывать некорректные записи как нулевые, попробуйте установить для параметра «mode» значение «PERMISSIVE».```
Как это исправить? моя конфигурация
public static void main(String[] args) throws StreamingQueryException, IOException, TimeoutException, RestClientException {
Map kafkaParams = new HashMap();
kafkaParams.put("bootstrap.servers", "localhost:9092");
kafkaParams.put("key.deserializer", AvroDeserializer.class);
kafkaParams.put("value.deserializer", AvroDeserializer.class);
kafkaParams.put("group.id", "sparkGroupId");
kafkaParams.put("auto.offset.reset", "earliest");
kafkaParams.put("enable.auto.commit", false);
SparkSession spark = SparkSession.builder().appName("Hello Spark").master(
"local").getOrCreate();
System.out.println("Hello, Spark v." + spark.version());
spark.sparkContext().setLogLevel("ERROR");
Subscribe to 1 topic
Dataset df = spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "j_client")
.option("startingOffsets", "earliest")
.load();
String schemaRegistryAddr = "http://localhost:8081";
создать объект Rest Service var restService = new RestService(schemaRegistryAddr);
// Create REST service to access schema registry and retrieve topic schema (latest)
var valueRestResponseSchema = restService.getLatestVersion("j_client" + "-value");
var jsonSchema = valueRestResponseSchema.getSchema();
var jClientDF =df.select(
col("key").cast("string"), // cast to string from binary value
from_avro(col("value"), jsonSchema).as("j_client"), // convert from avro value
col("topic"),
col("offset"),
col("timestamp"),
col("timestampType"))
;
jClientDF.printSchema();
Stream data to console for testing
jClientDF
.writeStream()
.format("console")
.start()
.awaitTermination()
;
Подробнее здесь: https://stackoverflow.com/questions/784 ... rd-parsing
Spark получает SparkException: при анализе записей обнаружены некорректные записи ⇐ JAVA
Программисты JAVA общаются здесь
1714591041
Anonymous
Я пытаюсь отобразить данные Avro из Kafka на консоли с помощью DataStream и поймать исключение.
Причина: org.apache.spark.SparkException: при анализе записей обнаружены некорректные записи. Текущий режим анализа: FAILFAST. Чтобы обрабатывать некорректные записи как нулевые, попробуйте установить для параметра «mode» значение «PERMISSIVE».```
Как это исправить? моя конфигурация
public static void main(String[] args) throws StreamingQueryException, IOException, TimeoutException, RestClientException {
Map kafkaParams = new HashMap();
kafkaParams.put("bootstrap.servers", "localhost:9092");
kafkaParams.put("key.deserializer", AvroDeserializer.class);
kafkaParams.put("value.deserializer", AvroDeserializer.class);
kafkaParams.put("group.id", "sparkGroupId");
kafkaParams.put("auto.offset.reset", "earliest");
kafkaParams.put("enable.auto.commit", false);
SparkSession spark = SparkSession.builder().appName("Hello Spark").master(
"local").getOrCreate();
System.out.println("Hello, Spark v." + spark.version());
spark.sparkContext().setLogLevel("ERROR");
Subscribe to 1 topic
Dataset df = spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "j_client")
.option("startingOffsets", "earliest")
.load();
String schemaRegistryAddr = "http://localhost:8081";
создать объект Rest Service var restService = new RestService(schemaRegistryAddr);
// Create REST service to access schema registry and retrieve topic schema (latest)
var valueRestResponseSchema = restService.getLatestVersion("j_client" + "-value");
var jsonSchema = valueRestResponseSchema.getSchema();
var jClientDF =df.select(
col("key").cast("string"), // cast to string from binary value
from_avro(col("value"), jsonSchema).as("j_client"), // convert from avro value
col("topic"),
col("offset"),
col("timestamp"),
col("timestampType"))
;
jClientDF.printSchema();
Stream data to console for testing
jClientDF
.writeStream()
.format("console")
.start()
.awaitTermination()
;
Подробнее здесь: [url]https://stackoverflow.com/questions/78415547/spark-getting-sparkexception-malformed-records-are-detected-in-record-parsing[/url]
Ответить
1 сообщение
• Страница 1 из 1
Перейти
- Кемерово-IT
- ↳ Javascript
- ↳ C#
- ↳ JAVA
- ↳ Elasticsearch aggregation
- ↳ Python
- ↳ Php
- ↳ Android
- ↳ Html
- ↳ Jquery
- ↳ C++
- ↳ IOS
- ↳ CSS
- ↳ Excel
- ↳ Linux
- ↳ Apache
- ↳ MySql
- Детский мир
- Для души
- ↳ Музыкальные инструменты даром
- ↳ Печатная продукция даром
- Внешняя красота и здоровье
- ↳ Одежда и обувь для взрослых даром
- ↳ Товары для здоровья
- ↳ Физкультура и спорт
- Техника - даром!
- ↳ Автомобилистам
- ↳ Компьютерная техника
- ↳ Плиты: газовые и электрические
- ↳ Холодильники
- ↳ Стиральные машины
- ↳ Телевизоры
- ↳ Телефоны, смартфоны, плашеты
- ↳ Швейные машинки
- ↳ Прочая электроника и техника
- ↳ Фототехника
- Ремонт и интерьер
- ↳ Стройматериалы, инструмент
- ↳ Мебель и предметы интерьера даром
- ↳ Cантехника
- Другие темы
- ↳ Разное даром
- ↳ Давай меняться!
- ↳ Отдам\возьму за копеечку
- ↳ Работа и подработка в Кемерове
- ↳ Давай с тобой поговорим...
Мобильная версия