This message was deleted.
# general
s
This message was deleted.
v
which version of druid are you on?
l
druid version is 25
v
please paste your ingestion spec here
l
Copy code
{
  "type": "kafka",
  "spec": {
    "dataSchema": {
      "dataSource": "test-duba",
      "timestampSpec": {
        "column": "content.TIMESTAMP",
        "format": "millis",
        "missingValue": null
      },
      "dimensionsSpec": {
        "dimensions": [
          {
            "type": "string",
            "name": "deviceType",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": true
          },
          {
            "type": "string",
            "name": "cseID",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": true
          },
          {
            "type": "string",
            "name": "_id",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": true
          },
          {
            "type": "string",
            "name": "fieldGateway",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": true
          },
          {
            "type": "string",
            "name": "deviceID",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": true
          },
          {
            "type": "string",
            "name": "content.TIMESTAMP",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": false
          },
          {
            "type": "string",
            "name": "content.rawValue",
            "multiValueHandling": "SORTED_ARRAY",
            "createBitmapIndex": false
          },
          {
            "type": "string",
            "name": "content.value",
            "multiValueHandling": "ARRAY",
            "createBitmapIndex": false
          }
        ],
        "dimensionExclusions": [
          "__time",
          "creationTime",
          "timestamp"
        ],
        "includeAllDimensions": false
      },
      "metricsSpec": [],
      "granularitySpec": {
        "type": "uniform",
        "segmentGranularity": "HOUR",
        "queryGranularity": {
          "type": "none"
        },
        "rollup": false,
        "intervals": []
      },
      "transformSpec": {
        "filter": null,
        "transforms": []
      }
    },
    "ioConfig": {
      "topic": "test-duba",
      "inputFormat": {
        "type": "avro_stream",
        "flattenSpec": {
          "useFieldDiscovery": true,
          "fields": [
            {
              "type": "path",
              "name": "content.TIMESTAMP",
              "expr": "$.content.TIMESTAMP",
              "nodes": null
            },
            {
              "type": "path",
              "name": "content.rawValue",
              "expr": "$.content.rawValue",
              "nodes": null
            },
            {
              "type": "path",
              "name": "content.value",
              "expr": "$.content.value",
              "nodes": null
            }
          ]
        },
        "avroBytesDecoder": {
          "type": "schema_inline",
          "schema": {
            "name": "AvroTelemetryDubaPlantScada",
            "type": "record",
            "namespace": "com.neos.dbconnector.model.dubaplantscada",
            "fields": [
              {
                "name": "_id",
                "type": "string"
              },
              {
                "name": "deviceID",
                "type": "string"
              },
              {
                "name": "deviceType",
                "type": "string"
              },
              {
                "name": "fieldGateway",
                "type": "string"
              },
              {
                "name": "cseID",
                "type": "string"
              },
              {
                "name": "creationTime",
                "type": "long"
              },
              {
                "name": "content",
                "type": {
                  "name": "AvroTelemetryDubaPlantScadaContent",
                  "type": "record",
                  "fields": [
                    {
                      "name": "rawValue",
                      "type": [
                        "null",
                        "double"
                      ]
                    },
                    {
                      "name": "TIMESTAMP",
                      "type": "long"
                    },
                    {
                      "name": "value",
                      "type": [
                        "null",
                        "string"
                      ]
                    }
                  ]
                }
              }
            ]
          }
        },
        "binaryAsString": false,
        "extractUnionsByType": false
      },
      "replicas": 1,
      "taskCount": 1,
      "taskDuration": "PT3600S",
      "consumerProperties": {
        "security.protocol": "SSL",
        "ssl.truststore.location": "***",
        "ssl.truststore.password": "***",
        "ssl.keystore.location": "***",
        "ssl.keystore.password": "***",
        "pollTimeout": 1000,
        "bootstrap.servers": "devkafka.core-dev.NEOS.today:9093"
      },
      "autoScalerConfig": null,
      "pollTimeout": 100,
      "startDelay": "PT5S",
      "period": "PT30S",
      "useEarliestOffset": true,
      "completionTimeout": "PT1800S",
      "lateMessageRejectionPeriod": null,
      "earlyMessageRejectionPeriod": null,
      "lateMessageRejectionStartDateTime": null,
      "configOverrides": null,
      "idleConfig": null,
      "stream": "test-duba",
      "useEarliestSequenceNumber": true,
      "type": "kafka"
    },
    "tuningConfig": {
      "type": "kafka",
      "appendableIndexSpec": {
        "type": "onheap",
        "preserveExistingMetrics": false
      },
      "maxRowsInMemory": 1000000,
      "maxBytesInMemory": 0,
      "skipBytesInMemoryOverheadCheck": false,
      "maxRowsPerSegment": 5000000,
      "maxTotalRows": null,
      "intermediatePersistPeriod": "PT10M",
      "maxPendingPersists": 0,
      "indexSpec": {
        "bitmap": {
          "type": "roaring",
          "compressRunOnSerialization": true
        },
        "dimensionCompression": "lz4",
        "stringDictionaryEncoding": {
          "type": "utf8"
        },
        "metricCompression": "lz4",
        "longEncoding": "longs"
      },
      "indexSpecForIntermediatePersists": {
        "bitmap": {
          "type": "roaring",
          "compressRunOnSerialization": true
        },
        "dimensionCompression": "lz4",
        "stringDictionaryEncoding": {
          "type": "utf8"
        },
        "metricCompression": "lz4",
        "longEncoding": "longs"
      },
      "reportParseExceptions": false,
      "handoffConditionTimeout": 0,
      "resetOffsetAutomatically": false,
      "segmentWriteOutMediumFactory": null,
      "workerThreads": null,
      "chatThreads": null,
      "chatRetries": 8,
      "httpTimeout": "PT10S",
      "shutdownTimeout": "PT80S",
      "offsetFetchPeriod": "PT30S",
      "intermediateHandoffPeriod": "P2147483647D",
      "logParseExceptions": false,
      "maxParseExceptions": 2147483647,
      "maxSavedParseExceptions": 0,
      "skipSequenceNumberAvailabilityCheck": false,
      "repartitionTransitionDuration": "PT120S"
    }
  },
  "context": null
}
v
@Lalit Verma is there any chance you can share a test avro file that replicates this issue?
l
the data is any grabbled data or corrupted data for testing I am sending an simple grabbled string 'asdfasdfasdf'
g
The configs for kafka are
maxParseExceptions
and
reportParseExceptions
(older, deprecated config)
It looks like you have already set them correctly: you did
reportParseExceptions: false
and
maxParseExceptions: 2147483647
(~2 billion, hopefully you did not hit this!)
I'm wondering what error you see exactly?
v
@Lalit Verma @Gian Merlino I took the spec and ran a task with it with just some random data in kafka and the task failed with
Copy code
org.apache.druid.java.util.common.parsers.ParseException: Avro's unnecessary EOFException, detail: [<https://issues.apache.org/jira/browse/AVRO-813>]
@Lalit Verma is this the error you are seeing?
g
@Vijay Narayanan do you have a full stack trace for this error?
v
Copy code
2023-03-01T03:52:23,217 INFO [task-runner-0-priority-0] org.apache.kafka.clients.Metadata - [Consumer clientId=consumer-kafka-supervisor-gjmihhhb-1, groupId=kafka-supervisor-gjmihhhb] Cluster ID: KeurxA5sT_eKFHuIHquSHg
2023-03-01T03:52:23,259 ERROR [task-runner-0-priority-0] org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner - Encountered exception in run() before persisting.
org.apache.druid.java.util.common.parsers.ParseException: Avro's unnecessary EOFException, detail: [<https://issues.apache.org/jira/browse/AVRO-813>]
	at org.apache.druid.data.input.avro.InlineSchemaAvroBytesDecoder.parse(InlineSchemaAvroBytesDecoder.java:92) ~[?:?]
	at org.apache.druid.data.input.avro.AvroStreamReader.intermediateRowIterator(AvroStreamReader.java:70) ~[?:?]
	at org.apache.druid.data.input.IntermediateRowParsingReader.intermediateRowIteratorWithMetadata(IntermediateRowParsingReader.java:231) ~[druid-core-25.0.0.jar:25.0.0]
	at org.apache.druid.data.input.IntermediateRowParsingReader.read(IntermediateRowParsingReader.java:49) ~[druid-core-25.0.0.jar:25.0.0]
	at org.apache.druid.segment.transform.TransformingInputEntityReader.read(TransformingInputEntityReader.java:43) ~[druid-processing-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.seekablestream.SettableByteEntityReader.read(SettableByteEntityReader.java:70) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.seekablestream.StreamChunkParser.parseWithInputFormat(StreamChunkParser.java:135) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.seekablestream.StreamChunkParser.parse(StreamChunkParser.java:104) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner.runInternal(SeekableStreamIndexTaskRunner.java:635) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner.run(SeekableStreamIndexTaskRunner.java:266) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.seekablestream.SeekableStreamIndexTask.runTask(SeekableStreamIndexTask.java:151) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.common.task.AbstractTask.run(AbstractTask.java:169) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.overlord.SingleTaskBackgroundRunner$SingleTaskBackgroundRunnerCallable.call(SingleTaskBackgroundRunner.java:477) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at org.apache.druid.indexing.overlord.SingleTaskBackgroundRunner$SingleTaskBackgroundRunnerCallable.call(SingleTaskBackgroundRunner.java:449) ~[druid-indexing-service-25.0.0.jar:25.0.0]
	at java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[?:?]
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) ~[?:?]
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) ~[?:?]
	at java.lang.Thread.run(Thread.java:834) ~[?:?]
Caused by: java.io.EOFException
	at org.apache.avro.io.BinaryDecoder$InputStreamByteSource.readRaw(BinaryDecoder.java:850) ~[?:?]
	at org.apache.avro.io.BinaryDecoder.doReadBytes(BinaryDecoder.java:372) ~[?:?]
	at org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:289) ~[?:?]
	at org.apache.avro.io.ResolvingDecoder.readString(ResolvingDecoder.java:209) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:469) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:459) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:191) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:160) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.readField(GenericDatumReader.java:259) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:247) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:179) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:160) ~[?:?]
	at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:153) ~[?:?]
	at org.apache.druid.data.input.avro.InlineSchemaAvroBytesDecoder.parse(InlineSchemaAvroBytesDecoder.java:88) ~[?:?]
	... 17 more
2023-03-01T03:52:23,268 INFO [task-runner-0-priority-0] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted rows[0] and (estimated) bytes[0]
2023-03-01T03:52:23,272 INFO [[index_kafka_test-duba_77ed6ecdfe1baa6_enmkjpdg]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Flushed in-memory data with commit metadata [AppenderatorDriverMetadata{segments={}, lastSegmentIds={}, callerMetadata={nextPartitions=SeekableStreamEndSequenceNumbers{stream='test-duba', partitionSequenceNumberMap={0=0}}}}] for segments: 
2023-03-01T03:52:23,272 INFO [[index_kafka_test-duba_77ed6ecdfe1baa6_enmkjpdg]-appenderator-persist] org.apache.druid.segment.realtime.appenderator.StreamAppenderator - Persisted stats: processed rows: [0], persisted rows[0], sinks: [0], total fireHydrants (across sinks): [0], persisted fireHydrants (across sinks): [0]
2023-03-01T03:52:23,272 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator - [Consumer clientId=consumer-kafka-supervisor-gjmihhhb-1, groupId=kafka-supervisor-gjmihhhb] Resetting generation and member id due to: consumer pro-actively leaving the group
2023-03-01T03:52:23,272 INFO [task-runner-0-priority-0] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator - [Consumer clientId=consumer-kafka-supervisor-gjmihhhb-1, groupId=kafka-supervisor-gjmihhhb] Request joining group due to: consumer pro-actively leaving the group
2023-03-01T03:52:23,273 INFO [task-runner-0-priority-0] org.apache.kafka.common.metrics.Metrics - Metrics scheduler closed
2023-03-01T03:52:23,273 INFO [task-runner-0-priority-0] org.apache.kafka.common.metrics.Metrics - Closing reporter org.apache.kafka.common.metrics.JmxReporter
l
@Vijay Narayanan @Gian Merlino that is the exact error which I am facing ... post which the state turns to be "UNHEALTHY_TASKS"
g
raised an issue with some more detail: https://github.com/apache/druid/issues/13894
I don't have a chance to fix it right now, but hopefully having an issue filed with the details helps someone else fix this!