Hi Team, I am trying to ingest data into DIMENSION...
# general
a
Hi Team, I am trying to ingest data into DIMENSION table with composite keys using spark-batch-ingestion pinot jar. But seems like ingestion is not complete i.e. There are rows in source which are not available in Pinot DIMENSION table.
Here are my schema, table and batch ingestion config. Schema:
Copy code
{
  "schemaName": "geohashAreaMapDim",
  "dimensionFieldSpecs": [
    {
      "name": "geohash",
      "dataType": "STRING"
    },
    {
      "name": "area",
      "dataType": "STRING"
    },
    {
      "name": "cityId",
      "dataType": "LONG"
    },
    {
      "name": "countryId",
      "dataType": "LONG"
    }
  ],
  "primaryKeyColumns": [
    "cityId",
    "geohash"
  ]
}
Table:
Copy code
{
  "OFFLINE": {
    "tableName": "geohashAreaMapDimOffline_OFFLINE",
    "tableType": "OFFLINE",
    "segmentsConfig": {
      "schemaName": "geohashAreaMapDim",
      "replication": "1",
      "segmentPushType": "REFRESH",
      "minimizeDataMovement": false
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant"
    },
    "tableIndexConfig": {
      "invertedIndexColumns": [],
      "rangeIndexVersion": 2,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "loadMode": "MMAP",
      "enableDefaultStarTree": false,
      "enableDynamicStarTreeCreation": false,
      "aggregateMetrics": false,
      "nullHandlingEnabled": false,
      "optimizeDictionaryForMetrics": false,
      "noDictionarySizeRatioThreshold": 0
    },
    "metadata": {},
    "quota": {
      "storage": "200M"
    },
    "ingestionConfig": {
      "batchIngestionConfig": {
        "segmentIngestionType": "REFRESH",
        "segmentIngestionFrequency": "DAILY"
      },
      "transformConfigs": [
        {
          "columnName": "cityId",
          "transformFunction": "city_id"
        },
        {
          "columnName": "countryId",
          "transformFunction": "country_id"
        }
      ]
    },
    "isDimTable": true
  }
}
Batch Ingestion Config:
Copy code
executionFrameworkSpec:
  name: 'spark'

  segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentGenerationJobRunner'
  segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentTarPushJobRunner'
  segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentUriPushJobRunner'
  extraConfigs:
    stagingDir: '<s3://bucket-name/midas-offline-staging/dimension/geohash.area_map>'

jobType: SegmentCreationAndTarPush
inputDirURI: '<s3://bucket-name/midas-offline-parquet/dimension/geohash.area_map/>'
includeFileNamePattern: 'glob:**/*.parquet'
outputDirURI: '<s3://bucket-name/midas-offline/dimension/geohash.area_map/>' # '/pinot/temp'
#auth-token
authToken: '<something>'

# overwriteOutput: Overwrite output segments if existed.
overwriteOutput: true

pinotFSSpecs:
  - scheme: s3
    className: org.apache.pinot.plugin.filesystem.S3PinotFS
    configs:
      region: 'ap-southeast-1'

recordReaderSpec:
  dataFormat: 'parquet'
  className: 'org.apache.pinot.plugin.inputformat.parquet.ParquetNativeRecordReader'

tableSpec:
  tableName: 'geohashAreaMapDimOffline'
  schemaURI: '<http://stg-mimic-pinot.coban.stg.com:9000/tables/geohashAreaMapDimOffline/schema>'
  tableConfigURI: '<http://stg-mimic-pinot.coban.stg.com:9000/tables/geohashAreaMapDimOffline>'

segmentNameGeneratorSpec:
  type: normalizedDate
  configs:
    segment.name.prefix: 'geohashAreaMapDimOffline_batch'
    exclude.sequence.id: true

pinotClusterSpecs:
  - controllerURI: '<http://stg-mimic-pinot.coban.stg.com:9000>'

pushJobSpec:
  pushParallelism: 2
  pushAttempts: 2
  pushRetryIntervalMillis: 1000
Pinot Output
Source Data Output:
As can be seen, there are 2 entries in source data for composity key (city_id, geohash) (106, w7db1n) (115, w7db1n) while in Pinot only 1 entry is available (city_id, geohash) (115, w7db1n)
seems like - primaryKeyColumns defined in schema isn't working as expected?
s
@Ashish Kumar how many files did you have for the input?
Can you double check if you have schema named as
geohashAreaMapDimOffline
?
Our code first tries to pick up the schema with the same name as the table name and then it will try to pull schema with the custom name. Also, the custom schema name is deprecated.
Copy code
1. change schema geohashAreaMapDim -> geohashAreaMapDimOffline
2. change the "schema" field in table config to geohashAreaMapDimOffline
a
@Seunghyun
how many files did you have for the input?
2 parquet files in the input.
Can you double check if you have schema named as
geohashAreaMapDimOffline
No, geohashAreaMapDim exists
s
Can you double check if 2 segments got ingested to pinot?
Screen Shot 2022-12-20 at 8.46.08 PM.png
a
no it's 1 segment only
s
it seems that the second file somehow did not get ingested?
we have 1-1 mapping from input file to a segment, How do you push data? Can you double check if the output segments from the different input didn't end up having the same segment name?
a
could it be because of my segment name generator spec:
Copy code
segmentNameGeneratorSpec:
  type: normalizedDate
  configs:
    segment.name.prefix: 'geohashAreaMapDimOffline_batch'
    exclude.sequence.id: true
s
Copy code
exclude.sequence.id: true -> false
can you retry?
a
okay, but if I add sequence id.. (sequence 0 & 1) and let's say in input one of the file is deleted (sequence 1).. then only sequence 0 is remaining then If i run the batch job again, will it also delete the segment segment corresponding to sequence 1 ?
s
if you run your batch job on the exactly the same data, it will create the same segment name with 0,1 and repush. In that case, 0 will be replaced and 1 will be newly added and we will end up having 0,1. So, each push job will be idempotent if the input data is the same
a
yeah that's okay, but I am saying what if input data changes i.e. no. of input files in same input folders are reduced. Running the batch ingestion job again on same input data folder will not exactly refresh the segment, right? better explained here https://apache-pinot.slack.com/archives/CDRCA57FC/p1652344365276709
s
@Ashish Kumar this has been the known issue for pinot. In the future, we will support this with the consistent push protocol. That will allow us to swap
old_0, old_1, old_2 -> new_0, new_1
. But this has not been added yet because the new protocol will now assume the time bucket based replacement (i.e. replace data for 1day at a time). Do you expect to backfill data frequently and the number of files for the original data can change frequently?
a
Do you expect to backfill data frequently and the number of files for the original data can change frequently?
yes, for now I am thinking of automating it using Pinot REST APIs... something like for backfills/original data changes.. (1) Delete all the existing segments with prefix! (2) Run the ingestion which will create new segments with new Data