Slackbot
01/09/2024, 10:42 AMLaksh Singla
01/09/2024, 11:21 AMLaksh Singla
01/09/2024, 11:21 AMTomasz Kisielewski
01/09/2024, 11:33 AMTomasz Kisielewski
01/09/2024, 11:36 AMMike Sherman
01/09/2024, 1:49 PMTomasz Kisielewski
01/09/2024, 2:40 PMMike Sherman
01/09/2024, 3:16 PMMike Sherman
01/09/2024, 3:17 PMTomasz Kisielewski
01/09/2024, 4:21 PMMike Sherman
01/09/2024, 4:58 PMLaksh Singla
01/09/2024, 6:36 PMI want to delete certain rows from datasource. In my case sometimes I need to change source files and then reingest into Druid, then, I have duplicates which I want to remove later onSounds like you’d wanna reingest the data source from an existing data source
modify the data that’s being re-ingestedCompaction is generally termed for reindexing with different clustering keys, rollup granularity etc. Modifying certain rows, or removing duplicates seem more like reingesting from a preexisting datasource. Both are same for something like “REPLACE”, but if going via native compaction, then you cannot do latter stuff (I think)
Laksh Singla
01/09/2024, 6:37 PMTomasz Kisielewski
01/09/2024, 10:05 PMTomasz Kisielewski
01/11/2024, 11:43 AM{"error":"Missing type id when trying to resolve subtype of [simple type, class org.apache.druid.query.search.SearchQuerySpec]: missing type id property 'type' (for POJO property 'query')\n at [Source: (org.eclipse.jetty.server.HttpInputOverHTTP); line: 1, column: 155832] (through reference chain: org.apache.druid.indexing.common.task.batch.parallel.ParallelIndexSupervisorTask[\"spec\"]->org.apache.druid.indexing.common.task.batch.parallel.ParallelIndexIngestionSpec[\"dataSchema\"]->org.apache.druid.segment.indexing.DataSchema[\"transformSpec\"
{
"type": "index_parallel",
"spec": {
"dataSchema": {
"dataSource": "test_data_2",
"timestampSpec": {
"column": "__time",
"format": "iso",
"missingValue": null
},
"dimensionsSpec": {
"includeAllDimensions": true,
"useSchemaDiscovery": false
},
"metricsSpec": [],
"granularitySpec": {
"type": "uniform",
"segmentGranularity": "DAY",
"queryGranularity": {
"type": "none"
},
"rollup": false,
"intervals": []
},
"transformSpec": {
"filter": {
"queryType": "scan",
"dataSource": {
"type": "join",
"left": {
"type": "table",
"name": "test_data"
},
"right": {
"type": "query",
"query": {
"queryType": "groupBy",
"dataSource": {
"type": "table",
"name": "test_data"
},
"intervals": {
"type": "intervals",
"intervals": [
"2023-10-01T00:00:00.000Z/2023-11-05T00:00:00.000Z"
]
},
"virtualColumns": [
{
"type": "expression",
"name": "v0",
"expression": "timestamp_parse(\"druid_ingestion_time\",null,\"UTC\")",
"outputType": "LONG"
}
],
"granularity": {
"type": "all"
},
"dimensions": [
{
"type": "default",
"dimension": "testid",
"outputName": "d0",
"outputType": "STRING"
}
],
"aggregations": [
{
"type": "longMax",
"name": "a0",
"fieldName": "v0"
}
],
"limitSpec": {
"type": "NoopLimitSpec"
}
}
},
"rightPrefix": "j0.",
"condition": "((\"testid\" == \"j0.d0\") && (timestamp_parse(\"druid_ingestion_time\",null,\"UTC\") == \"j0.a0\"))",
"joinType": "LEFT"
},
"intervals": {
"type": "intervals",
"intervals": [
"2023-10-01T00:00:00.000Z/2023-11-05T00:00:00.000Z"
]
},
"virtualColumns": [
{
"type": "expression",
"name": "v0",
"expression": "nvl(\"j0.d0\",\"null\")",
"outputType": "STRING"
}
],
"resultFormat": "compactedList",
"filter": {
"type": "not",
"field": {
"type": "selector",
"dimension": "v0",
"value": "null"
}
},
"columns": [
"testid",
"subtestname",
"druid_ingestion_time"
],
"legacy": false,
"granularity": {
"type": "all"
}
}
}
},
"ioConfig": {
"type": "index_parallel",
"inputSource": {
"type": "druid",
"dataSource": "test_data",
"interval": "2023-10-01/2023-11-05"
}
},
"tuningConfig": {
"type": "index_parallel",
"maxRowsPerSegment": 15000000,
"appendableIndexSpec": {
"type": "onheap",
"preserveExistingMetrics": false
},
"maxRowsInMemory": 50000000,
"maxBytesInMemory": 0,
"skipBytesInMemoryOverheadCheck": false,
"maxTotalRows": 20000000,
"numShards": null,
"splitHintSpec": {
"type": "maxSize",
"maxSplitSize": 104857600,
"maxNumFiles": 1000
},
"partitionsSpec": {
"type": "dynamic",
"maxRowsPerSegment": 15000000,
"maxTotalRows": 20000000
},
"indexSpec": {
"bitmap": {
"type": "roaring"
},
"dimensionCompression": "lz4",
"stringDictionaryEncoding": {
"type": "utf8"
},
"metricCompression": "lz4",
"longEncoding": "longs"
},
"indexSpecForIntermediatePersists": {
"bitmap": {
"type": "roaring"
},
"dimensionCompression": "lz4",
"stringDictionaryEncoding": {
"type": "utf8"
},
"metricCompression": "lz4",
"longEncoding": "longs"
},
"maxPendingPersists": 0,
"forceGuaranteedRollup": false,
"reportParseExceptions": false,
"pushTimeout": 0,
"segmentWriteOutMediumFactory": null,
"maxNumConcurrentSubTasks": 58,
"maxRetry": 3,
"taskStatusCheckPeriodMs": 1000,
"chatHandlerTimeout": "PT10S",
"chatHandlerNumRetries": 5,
"maxNumSegmentsToMerge": 100,
"totalNumMergeTasks": 10,
"logParseExceptions": true,
"maxParseExceptions": 2147483647,
"maxSavedParseExceptions": 0,
"maxColumnsToMerge": -1,
"awaitSegmentAvailabilityTimeoutMillis": 0,
"maxAllowedLockCount": -1,
"partitionDimensions": []
}
},
"context": {
"forceTimeChunkLock": true,
"useLineageBasedSegmentAllocation": true
}
}
which translates to:
select * from (
with maxx as (
select
distinct testid, max(cast(druid_ingestion_time as timestamp)) druid_ingestion_time
from "test_data"
group by testid
)
select
a.*, COALESCE(maxx.testid, 'null') as testid2
from "test_data" a
left join maxx on maxx.testid = a.testid and maxx.druid_ingestion_time = cast(a.druid_ingestion_time as timestamp)
)
where __time >= TIMESTAMP '2023-10-01 00:00:00' and __time < TIMESTAMP '2023-11-05 00:00:00'
and testid2 <> 'null'Laksh Singla
01/11/2024, 6:37 PM