Hi team, I have a table with 13 primary keys and <...
# general
a
Hi team, I have a table with 13 primary keys and upsert enabled. Data volume is 1k messages/second and 900KB/second. I have distributed pinot traffic to 10 servers. Does anyone know what’s the suggested heap size for a pinot server? Thanks. Table config is as below:
Copy code
{
  "REALTIME": {
    "tableName": "howler_ad_mainst_battlestation_order_updates_REALTIME",
    "tableType": "REALTIME",
    "segmentsConfig": {
      "schemaName": "howler_ad_mainst_battlestation_order_updates",
      "replication": "1",
      "replicasPerPartition": "1",
      "timeColumnName": "time",
      "minimizeDataMovement": false
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant",
      "tagOverrideConfig": {}
    },
    "tableIndexConfig": {
      "invertedIndexColumns": [],
      "noDictionaryColumns": [],
      "streamConfigs": {
        "streamType": "kafka",
        "stream.kafka.topic.name": "howler_ad_mainst_battlestation_order_updates",
        "stream.kafka.broker.list": "confluent-broker.roles.service.robinhood:9092",
        "stream.kafka.consumer.type": "lowlevel",
        "stream.kafka.consumer.prop.auto.offset.reset": "smallest",
        "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory",
        "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
        "realtime.segment.flush.threshold.rows": "0",
        "realtime.segment.flush.threshold.time": "24h",
        "realtime.segment.flush.segment.size": "100M"
      },
      "rangeIndexColumns": [],
      "rangeIndexVersion": 2,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "sortedColumn": [],
      "bloomFilterColumns": [],
      "loadMode": "MMAP",
      "onHeapDictionaryColumns": [],
      "varLengthDictionaryColumns": [],
      "enableDefaultStarTree": false,
      "enableDynamicStarTreeCreation": false,
      "aggregateMetrics": false,
      "nullHandlingEnabled": false,
      "optimizeDictionaryForMetrics": false,
      "noDictionarySizeRatioThreshold": 0
    },
    "metadata": {},
    "quota": {},
    "routing": {
      "instanceSelectorType": "strictReplicaGroup"
    },
    "query": {},
    "upsertConfig": {
      "mode": "FULL",
      "hashFunction": "NONE"
    },
    "ingestionConfig": {},
    "isDimTable": false
  }
}
m
You have 13 primary keys? Is this is de-dup use case or upsert use case?
a
Either dedup or upsert works. In my understanding, if there are multiple primary keys which have the same values, dedup will store the first event while upsert will store the last event, right?
m
Yes
a
Which option do you suggest? Does dedup consume more memory than upsert?
m
You probably want to use upsert (more tested than dedup). Also, you want to probably use hash-code for PK (13 columns will take up a lot of memory)
a
Got it. And our kafka topics has partitions so Pinot table which ingests data from it have segments in different servers. If the kafka source event are not sent with the key of 13 primary columns, will dedup or upsert fail?
m
Yes, PK’s are expected to be always present.
a
Last question is if we have hash code for PK, how much memory do you estimate we should set for 1 pinot server?
m
To tag along in the discussion. If we partition our Flink stream does our keying method matter? Can we just run a hash over the output string or do we have to run the same hash algorithm/columns that Pinot is? The page says it needs to be partitioned by the primary key, but does this require the key to be the same as Pinot's hash algorithm?
a
My understanding is they don’t need to be the same. Correct me if I am wrong @Mayank : The reason we need to partition the input kafka msgs is to make sure data with the same primary key values are stored in the same pinot server. Later Pinot can do dedup by its own hash functions. Suppose kafka sending and pinot use different hash functions. Those msgs with the same primary key values are still stored in the same server by kafka sending hash value and Pinot is still able to dedup them by the Pinot hash value.
m
@Matthew Kerian partition function algorithm doesn’t need to match for upsert. But if you want to make sure query goes to only one server (performance optimization, unrelated to upsert), then yes needs to match
a
I don’t understand. How kafka’s and pinot’s hash functions match guarantee queries going to one server? Suppose I have multiple rows with different primary key value groups. Those groups are stored in different pinot servers. If I want to query SELECT COUNT(*), even though kafka hash code and pinot hash code are teh same, pinot still needs to go to different servers to get the full count, right?
m
Does Pinot hash longs big endian or little endian do you know?
g
cc @Mingmin Xu