i loss records in my kafka streams.
My kafka stream is in spark infra with that config.
val df = spark.
readStream.
format("kafka").
option("kafka.bootstrap.servers", broker_address).
option("subscribe", subject).
option("startingOffsets", "latest").
option("failOnDataLoss", "false").
load()
My sink is a parquet files.
df_vertex.writeStream.
format("parquet").
option("checkpointLocation", "/tmp/vertex/check").
option("path", data_location).
option("mode", "append").
trigger(Trigger.ProcessingTime("10 seconds")).
start().
awaitTermination()
When a read my parquets files some records are missing. That's recors are present in kafka broker but not write on sink.
21/12/04 05:32:50 DEBUG KafkaDataConsumer: Get spark-kafka-source-2edd8854-b64d-4f1a-b130-135e0c2b4c56--578079920-executor rainbow_data_extractor-0 nextOffset 42141 requested $offset
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: <!-- start message -->
21/12/04 05:32:50 DEBUG MessageColumnIO: < MESSAGE START >
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={}}: [] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: <id>
21/12/04 05:32:50 DEBUG MessageColumnIO: startField(id, 0)
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={}}: [id] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: [99, 102, 49, 48, 49, 48, 99, 98, 100, 53, 53, 49, 48, 51, 49, 98, 55, 53, 51, 97, 100, 55, 101, 49, 100, 99, 50, 49, 101, 100, 100, 100, 101, 51, 53, 55, 101, 100, 100, 49, 98, 99, 98, 51, 53, 101, 55, 100, 99, 97, 49, 56, 102, 51, 100, 49, 49, 55, 54, 53, 98, 5
5, 48, 56]
21/12/04 05:32:50 DEBUG MessageColumnIO: addBinary(64 bytes)
21/12/04 05:32:50 DEBUG MessageColumnIO: r: 0
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={}}: [id] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: </id>
21/12/04 05:32:50 DEBUG MessageColumnIO: endField(id, 0)
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={0}}: [] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: <type>
21/12/04 05:32:50 DEBUG MessageColumnIO: startField(type, 1)
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={0}}: [type] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: [117, 115, 101, 114]
21/12/04 05:32:50 DEBUG MessageColumnIO: addBinary(4 bytes)
21/12/04 05:32:50 DEBUG MessageColumnIO: r: 0
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={0}}: [type] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: </type>
21/12/04 05:32:50 DEBUG MessageColumnIO: endField(type, 1)
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={0, 1}}: [] r:0
21/12/04 05:32:50 DEBUG MessageColumnIO: [created_by].writeNull(0,0)
21/12/04 05:32:50 DEBUG MessageColumnIO: [created_on].writeNull(0,0)
21/12/04 05:32:50 DEBUG MessageColumnIO: [last_message].writeNull(0,0)
21/12/04 05:32:50 DEBUG MessageColumnIO: [is_active].writeNull(0,0)
21/12/04 05:32:50 DEBUG MessageColumnIO: [is_archived].writeNull(0,0)
21/12/04 05:32:50 DEBUG MessageColumnIO: < MESSAGE END >
21/12/04 05:32:50 DEBUG MessageColumnIO: 0, VistedIndex{vistedIndexes={0, 1}}: [] r:0
21/12/04 05:32:50 DEBUG RecordConsumerLoggingWrapper: <!-- end message -->
21/12/04 05:32:50 DEBUG KafkaDataConsumer: Get spark-kafka-source-2edd8854-b64d-4f1a-b130-135e0c2b4c56--578079920-executor rainbow_data_extractor-0 nextOffset 42142 requested $offset
21/12/04 05:32:50 DEBUG KafkaDataConsumer: Get spark-kafka-source-2edd8854-b64d-4f1a-b130-135e0c2b4c56--578079920-executor rainbow_data_extractor-0 nextOffset 42143 requested $offset
21/12/04 05:32:50 DEBUG KafkaDataConsumer: Get spark-kafka-source-2edd8854-b64d-4f1a-b130-135e0c2b4c56--578079920-executor rainbow_data_extractor-0 nextOffset 42144 requested $offset
21/12/04 05:32:50 DEBUG KafkaDataConsumer: Get spark-kafka-source-2edd8854-b64d-4f1a-b130-135e0c2b4c56--578079920-executor rainbow_data_extractor-0 nextOffset 42145 requested $offset
21/12/04 05:32:50 DEBUG KafkaDataConsumer: Get spark-kafka-source-2edd8854-b64d-4f1a-b130-135e0c2b4c56--578079920-executor rainbow_data_extractor-0 nextOffset 42146 requested $offset
On this example i loss records from 42142 to 42146. All my records are put in kafka broker in a short time. Perhaps i have some pb in a received scheduling of records.
If someone have the right configuration of stream for my case ?
Thanks Seb