This message was deleted.
# troubleshooting
s
This message was deleted.
k
n
getting this err when used in transformspec { "ingestionState": "BUILD_SEGMENTS", "unparseableEvents": {}, "rowStats": { "buildSegments": { "processed": 0, "processedWithError": 0, "thrownAway": 0, "unparseable": 0 } }, "errorMsg": "org.apache.druid.java.util.common.ISE: Could not transform value for kafka.timestamp reason: Function[timestamp] first argument should be a STRING but got LONG instead\n\tat org.apache.druid.segment.transform.ExpressionTransform$ExpressionRowFunction.eval(ExpressionTransform.java:117)\n\tat org.apache.druid.segment.transform.Transformer$TransformedInputRow.getRaw(Transformer.java:218)\n\tat org.apache.druid.segment.incremental.IncrementalIndex.toIncrementalIndexRow(IncrementalIndex.java:581)\n\tat org.apache.druid.segment.incremental.IncrementalIndex.add(IncrementalIndex.java:516)\n\tat org.apache.druid.segment.realtime.plumber.Sink.add(Sink.java:184)\n\tat org.apache.druid.segment.realtime.appenderator.StreamAppenderator.add(StreamAppenderator.java:276)\n\tat org.apache.druid.segment.realtime.appenderator.BaseAppenderatorDriver.append(BaseAppenderatorDriver.java:411)\n\tat org.apache.druid.segment.realtime.appenderator.StreamAppenderatorDriver.add(StreamAppenderatorDriver.java:189)\n\tat org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner.runInternal(SeekableStreamIndexTaskRunner.java:654)\n\tat org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner.run(SeekableStreamIndexTaskRunner.java:266)\n\tat org.apache.druid.indexing.seekablestream.SeekableStreamIndexTask.run(SeekableStreamIndexTask.java:151)\n\tat org.apache.druid.indexing.overlord.SingleTaskBackgroundRunner$SingleTaskBackgroundRunnerCallable.call(SingleTaskBackgroundRunner.java:477)\n\tat org.apache.druid.indexing.overlord.SingleTaskBackgroundRunner$SingleTaskBackgroundRunnerCallable.call(SingleTaskBackgroundRunner.java:449)\n\tat java.util.concurrent.FutureTask.run(FutureTask.java:266)\n\tat java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)\n\tat java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)\n\tat java.lang.Thread.run(Thread.java:750)\n\tSuppressed: java.lang.RuntimeException: java.lang.IllegalArgumentException: fromIndex(0) > toIndex(-1)\n\t\tat org.apache.druid.segment.realtime.appenderator.StreamAppenderatorDriver.persist(StreamAppenderatorDriver.java:241)\n\t\tat org.apache.druid.indexing.seekablestream.SeekableStreamIndexTaskRunner.runInternal(SeekableStreamIndexTaskRunner.java:772)\n\t\t... 8 more\n\tCaused by: java.lang.IllegalArgumentException: fromIndex(0) > toIndex(-1)\n\t\tat java.util.ArrayList.subListRangeCheck(ArrayList.java:1016)\n\t\tat java.util.ArrayList.subList(ArrayList.java:1006)\n\t\tat org.apache.druid.segment.realtime.appenderator.StreamAppenderator.persistAll(StreamAppenderator.java:589)\n\t\tat org.apache.druid.segment.realtime.appenderator.StreamAppenderatorDriver.persist(StreamAppenderatorDriver.java:233)\n\t\t... 9 more\nCaused by: org.apache.druid.math.expr.ExpressionValidationException: Function[timestamp] first argument should be a STRING but got LONG instead\n\tat org.apache.druid.math.expr.NamedFunction.validationFailed(NamedFunction.java:46)\n\tat org.apache.druid.math.expr.Function$TimestampFromEpochFunc.apply(Function.java:2829)\n\tat org.apache.druid.math.expr.FunctionExpr.eval(FunctionalExpr.java:183)\n\tat org.apache.druid.segment.transform.ExpressionTransform$ExpressionRowFunction.eval(ExpressionTransform.java:113)\n\t... 16 more\n", "segmentAvailabilityConfirmed": false, "segmentAvailabilityWaitTimeMs": 0 }
"transformSpec": { "filter": null, "transforms": [ { "type": "expression", "name": "kafka.timestamp", "expression": "timestamp(\"kafka.timestamp\")" } ] } },
From the log it says it is long, but expecting string , not sure how it is changed to LONG
g
I think it already comes out as long from Kafka?
You should be able to define it as a dimension of type
long
without using a transform
Does that work?
n
it is coming as string, is there any possibility to do it instead of adding in dimension. since we are not aware of all dimensions so we kept our ingestion spec like this "dimensionsSpec": { "dimensions": [], "dimensionExclusions": [ "__time", "ts" ], "includeAllDimensions": true } if we define only kafka.timestamp in dimension other dimensions are not seen in the query
g
in this case everything is loaded as a string. it's a limitation with the schemaless feature today
we are working on type detection so they will be loaded as the proper type, @Eric Tschetter or @Clint Wylie should know more
til then you may also experiment with the https://druid.apache.org/docs/latest/querying/nested-columns.html feature — not sure if it will satisfy your needs, but it also provides a way of loading data without knowing all the columns in advance, so it's worth a look
and it is possible to combine this with explicitly-specified dimensions, so you can get the correct type for your kafka timestamp
n
Can we do it using transformSpec, so it will be simple? A final option that we see is using a query, just if data is already converted to the required format we don't need conversions in query :)
we saw there is a parameter includeAllDimensions (https://druid.apache.org/docs/24.0.0/ingestion/ingestion-spec.html#dimensionsspec) which says "You can set
includeAllDimensions
it to true to ingest both explicit dimensions in the
dimensions
field and other dimensions that the ingestion task discovers from input data", tried setting this parameter but ended up only with kafka.timestamp column and other columns are missing. by enabling this are we going to have schemaless dimensions plus "kafka.timestamp" with required type ? tried the same and we see no changes, can u please the spec and let me know if it is correct?
Hi @Gian Merlino can you please check and help with this? are we missing something ?
e
What version are you running?
n
we are using 24.0.0
e
Hrm... From reading the code, that includeAllDimensions should work. When you say you are only getting kafka.timestamp and nothing else, is that from the sampling in the wizard in the console or are you actually running the ingestion and not getting the columns?
I ask because we sometimes see bugs where the console shows one thing due to something weird the Sampler is doing while the ingestion jobs themselves do the right thing. I'm curious if this is going to be one of those cases.
n
this our sample data: {"ts":"2021-11-29T174500+08:00","mt":"pwr","p1":"BLR","p2":"767","p3":"767","p4":"171","p5":"244","p6":"08","p7":"001","p8":"001","p9":"test","p12345":930,"p233434":900,"duration":15} when we ingest this data with the above spec in the SQL console we see only "kafka.timestamp"(so basically it is showing only 'kafka.timestamp' and 'ts' and other columns are not shown in the SQL query tab)
e
Oh, do you have the supervisor/task still running?
n
yes
e
Can you pause the supervisor?
That should cause it to end the task and that will cause a hand-off of the segment. I suspect that will fix the query console
"fix"
Also, make sure that you are refreshing your browser. The schema in the browser is grabbed once and not updated.
n
we paused(suspended) the supervisor and tried refreshing browser, but still, we see the same.
PFA images of same.
e
Can you go to the segments view, click on the magnifying glass to the right of the line for that one segment and copy&paste or screenshot the Payload section of that?
n
{ "dataSource": "test", "interval": "2021-11-29T090000.000Z/2021-11-29T100000.000Z", "version": "2022-10-31T052448.269Z", "loadSpec": { "type": "hdfs", "path": "hdfs://hadoop.local:8020/druid/segments/test/20211129T090000.000Z_20211129T100000.000Z/2022-10-31T05_24_48.269Z/0_1c5681a1-6443-4bad-81b9-718216fbfa7b_index.zip" }, "dimensions": "kafka.timestamp", "metrics": "", "shardSpec": { "type": "numbered", "partitionNum": 0, "partitions": 0 }, "binaryVersion": 9, "size": 739, "identifier": "test_2021-11-29T090000.000Z_2021-11-29T100000.000Z_2022-10-31T052448.269Z" }
is this the one are u referring to ?
e
That is the one... and hrm. That is weird. I'm not sure exactly what is going on, it's definitely not what I expect. If it's possible, you can use the nested column to do auto-discovered columns, which will also identify types as well
n
can u plz share some links or examples for this nested column?
n
will go through this , Thank you.
r
Hi Naresh, one way to do this is to pre-process the data. In other words, send the data from Kafka to a real time stream processor, convert the string to long, and deliver to Druid.