This message was deleted.
# general
s
This message was deleted.
l
Nope, unless you wanna rebuild Druid
Do you wanna delete the whole datasource, or just a portion of it ?
t
I 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 on
is native batch ingestion reindex my only option in that case?
m
Hey @Laksh Singla doesn't compaction allow you to modify the data that's being re-ingested? Not sure that's going to help in this case but the use case sounds like post-processing which compaction is designed to do.
t
@Mike Sherman, It would be very good but how, by modifying transform spec? Any example how to preserve only unique values based on ID and last_updated (max value) column?
m
Ya, @Tomasz Kisielewski I don't know the details as I no longer have access to that Druid cluster. In general, it's an anti-pattern to ingest duplicate rows into your forward cache. I wrote a small application at my last job where I read the Kafka events going into Druid and wrote them to HDFS as a backup, and to provide access outside of Druid. You might consider having some kind of application like that to clean up the Kafka events prior to publishing them to be consumed into your Druid datasource.
We handled "duplicate" rows by enabling rollup - https://druid.apache.org/docs/latest/tutorials/tutorial-rollup/. Of course, the drawback here is that your metrics get aggregated...
t
Thank you Mike, this is for sure useful, but not in my case 🙂 Neither I use Kafka (batch, and we need data ASAP, and changes occur a few days later) nor interested in aggregations. We assumed at the very beginning that the number of changed in source data will be small enough that we will handle them using proper filtering in query. Unfortunately, life is brutal 🙂 We have more changes that we want to and queries start to be more and more inefficient at scale. Looking for a solution and "replace into" was kind of trick, but 2000 columns? Are we in 80's? 😉
m
Ya, I feel for you in terms of the queries and filtering. We had some trillion row datasets at my last job but they weren't as wide as you're seeing. You might look into unioning the datasets at query time and breaking up the width, I've seen this used. If you query against a dimension or metric that's not in one of the union components, it's just ignored. Perhaps try a small test...
l
I 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 on
Sounds like you’d wanna reingest the data source from an existing data source
modify the data that’s being re-ingested
Compaction 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)
@Tomasz Kisielewski You can aggregate using the LATEST_BY aggregation. checkout the Druid docs for the latest function.
t
@Laksh Singla, you mean to use latest_by in queries, compaction or ingestion? I'm confused. I want to re-ingest data from existing datasource as you said, yet I do not know how if not using sql-based ingestion
Can I use "queryType": "scan" in filter of input from datasource? I wanted to do something like below but it throws me an error:
{"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\"
Copy code
{
  "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:
Copy code
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'
l
@Tomasz Kisielewski the native ingestion spec is different from the query that you are runninghttps://druid.apache.org/docs/latest/ingestion/ingestion-spec#transformspec