I am observing inconsistent results in my Pinot se...
# general
a
I am observing inconsistent results in my Pinot select Query.We are using upsert.Can someone help me with this
m
@Kartik Khare ^^
k
Hi, can you share the query and results here
r
@Abhijeet Kushe How is your Kafka topic partitioned?
a
@robert zych We are using Kinesis
Screen Shot 2023-01-20 at 6.41.30 AM.png,Screen Shot 2023-01-20 at 6.41.15 AM.png
@Kartik Khare this the same query but different results
Copy code
select ToDateTime(eventTimestamp, 'yyyy-MM-dd HH:mm:ss.SSS') as eventTimestamp,createdOn, $segmentName, "accountId", "workflowDefinitionId", "workflowInstanceId", "workflowRunningId", "recordType", "taskId", "taskKind", "attributeId", "automationFlowId", "contactId", "campaignActivityId", "sendId",orderNumber, softDelete  from workflowEvents where accountId = 1100678423876 AND campaignActivityId = '528680c5-6b72-4e7d-8a9b-febe061d28e7' AND activityType = 'null' AND recordType = 'attribution' AND softDelete = 'null' order by eventTimestamp desc limit 50
Copy code
eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	softDelete
2023-01-19 21:11:54.709	1642972105000	workflowEvents__1__191__20230119T1352Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:Shopify:1003	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	a9a0c958-39da-4d39-932e-118214590009	1003	null
Copy code
Output #2

eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	softDelete
2023-01-19 21:11:54.709	1642972105000	workflowEvents__1__191__20230119T1352Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:Shopify:1003	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	a9a0c958-39da-4d39-932e-118214590009	1003	null
2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__0__20220720T1533Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null
k
can you also add $hostName to the query and share the results again
a
Copy code
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	softDelete
cdp-dl-pinot-k8s-server-0.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2023-01-19 21:11:54.709	1642972105000	workflowEvents__1__191__20230119T1352Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:Shopify:1003	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	a9a0c958-39da-4d39-932e-118214590009	1003	null
Copy code
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	softDelete
cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2023-01-19 21:11:54.709	1642972105000	workflowEvents__1__191__20230119T1352Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:Shopify:1003	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	a9a0c958-39da-4d39-932e-118214590009	1003	null
cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__0__20220720T1533Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null
k
seems like the response of server-0 and server-2 are different if this is not a partial upsert table can you try reloading the segment
a
I did do a reload last Friday as well from the console .But the results did not change
I have tried it again
k
can you share the upsert config and schema as well if possible
a
Copy code
{
  "REALTIME": {
    "tableName": "workflowEvents_REALTIME",
    "tableType": "REALTIME",
    "segmentsConfig": {
      "timeType": "MILLISECONDS",
      "schemaName": "workflowEvents",
      "retentionTimeUnit": "DAYS",
      "retentionTimeValue": "1826",
      "timeColumnName": "eventTimestamp",
      "allowNullTimeValue": false,
      "replicasPerPartition": "3",
      "segmentPushType": "APPEND"
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant"
    },
    "tableIndexConfig": {
      "streamConfigs": {
        "streamType": "kinesis",
        "stream.kinesis.topic.name": "qa-events-stream",
        "region": "us-east-1",
        "shardIteratorType": "LATEST",
        "stream.kinesis.consumer.type": "lowlevel",
        "stream.kinesis.fetch.timeout.millis": "30000",
        "stream.kinesis.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
        "stream.kinesis.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kinesis.KinesisConsumerFactory",
        "realtime.segment.flush.threshold.size": "5000000",
        "realtime.segment.flush.threshold.time": "1d"
      },
      "rangeIndexVersion": 1,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "loadMode": "MMAP",
      "enableDefaultStarTree": false,
      "aggregateMetrics": false,
      "enableDynamicStarTreeCreation": false,
      "nullHandlingEnabled": false
    },
    "metadata": {
      "customConfigs": {}
    },
    "routing": {
      "instanceSelectorType": "strictReplicaGroup"
    },
    "upsertConfig": {
      "mode": "FULL",
      "hashFunction": "NONE"
    },
    "isDimTable": false
  }
}
Copy code
{
  "schemaName": "workflowEvents",
  "dimensionFieldSpecs": [
    {
      "name": "accountId",
      "dataType": "LONG"
    },
    {
      "name": "recordType",
      "dataType": "STRING"
    },
    {
      "name": "workflowDefinitionId",
      "dataType": "STRING"
    },
    {
      "name": "workflowDefinitionName",
      "dataType": "STRING"
    },
    {
      "name": "workflowDefinitionVersion",
      "dataType": "STRING"
    },
    {
      "name": "workflowInstanceId",
      "dataType": "STRING"
    },
    {
      "name": "workflowRunningId",
      "dataType": "STRING"
    },
    {
      "name": "workflowPrimaryId",
      "dataType": "STRING"
    },
    {
      "name": "workflowSecondaryId",
      "dataType": "STRING"
    },
    {
      "name": "workflowStatus",
      "dataType": "STRING"
    },
    {
      "name": "taskId",
      "dataType": "STRING"
    },
    {
      "name": "taskName",
      "dataType": "STRING"
    },
    {
      "name": "taskKind",
      "dataType": "STRING"
    },
    {
      "name": "taskStatus",
      "dataType": "STRING"
    },
    {
      "name": "taskResult",
      "dataType": "STRING"
    },
    {
      "name": "taskSkipped",
      "dataType": "STRING"
    },
    {
      "name": "campaignId",
      "dataType": "STRING"
    },
    {
      "name": "campaignActivityId",
      "dataType": "STRING"
    },
    {
      "name": "contactId",
      "dataType": "STRING"
    },
    {
      "name": "automationFlowId",
      "dataType": "STRING"
    },
    {
      "name": "automationTemplateId",
      "dataType": "STRING"
    },
    {
      "name": "sendId",
      "dataType": "STRING"
    },
    {
      "name": "currencyCode",
      "dataType": "STRING"
    },
    {
      "name": "attributeId",
      "dataType": "STRING"
    },
    {
      "name": "channel",
      "dataType": "STRING"
    },
    {
      "name": "correlationId",
      "dataType": "STRING"
    },
    {
      "name": "activityType",
      "dataType": "STRING"
    },
    {
      "name": "softDelete",
      "dataType": "STRING"
    },
    {
      "name": "workflowNamespace",
      "dataType": "STRING",
      "defaultNullValue": "ctct:dp:rmf"
    },
    {
      "name": "workflowDisplayName",
      "dataType": "STRING"
    },
    {
      "name": "workflowAccountId",
      "dataType": "STRING"
    },
    {
      "name": "attributeType",
      "dataType": "STRING"
    },
    {
      "name": "attributionType",
      "dataType": "STRING"
    },
    {
      "name": "storeFrontName",
      "dataType": "STRING"
    },
    {
      "name": "storeFrontType",
      "dataType": "STRING"
    },
    {
      "name": "checkoutToken",
      "dataType": "STRING"
    },
    {
      "name": "cartToken",
      "dataType": "STRING"
    },
    {
      "name": "orderNumber",
      "dataType": "STRING"
    },
    {
      "name": "bounceCode",
      "dataType": "STRING"
    },
    {
      "name": "urlId",
      "dataType": "STRING"
    },
    {
      "name": "linkUrl",
      "dataType": "STRING"
    },
    {
      "name": "messageId",
      "dataType": "STRING"
    },
    {
      "name": "isBilled",
      "dataType": "STRING"
    },
    {
      "name": "deviceType",
      "dataType": "STRING"
    }
  ],
  "metricFieldSpecs": [
    {
      "name": "sentCount",
      "dataType": "LONG"
    },
    {
      "name": "openCount",
      "dataType": "LONG"
    },
    {
      "name": "clickCount",
      "dataType": "LONG"
    },
    {
      "name": "bounceCount",
      "dataType": "LONG"
    },
    {
      "name": "totalOrderAmount",
      "dataType": "LONG"
    },
    {
      "name": "totalOrderAmt",
      "dataType": "DOUBLE"
    },
    {
      "name": "send",
      "dataType": "LONG"
    },
    {
      "name": "open",
      "dataType": "LONG"
    },
    {
      "name": "click",
      "dataType": "LONG"
    },
    {
      "name": "bounce",
      "dataType": "LONG"
    },
    {
      "name": "deliver",
      "dataType": "LONG"
    },
    {
      "name": "queue",
      "dataType": "LONG"
    },
    {
      "name": "totalOrderNum",
      "dataType": "LONG"
    },
    {
      "name": "discountAmount",
      "dataType": "DOUBLE"
    },
    {
      "name": "tax1Amount",
      "dataType": "DOUBLE"
    },
    {
      "name": "numberOfMessageParts",
      "dataType": "LONG"
    }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "eventTimestamp",
      "dataType": "LONG",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    },
    {
      "name": "eventTimestampNanos",
      "dataType": "LONG",
      "format": "1:NANOSECONDS:EPOCH",
      "granularity": "1:NANOSECONDS"
    },
    {
      "name": "createdOn",
      "dataType": "LONG",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    },
    {
      "name": "sendCreatedOn",
      "dataType": "LONG",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    },
    {
      "name": "orderLastModifiedDate",
      "dataType": "LONG",
      "format": "1:MILLISECONDS:EPOCH",
      "granularity": "1:MILLISECONDS"
    }
  ],
  "primaryKeyColumns": [
    "accountId",
    "workflowDefinitionId",
    "workflowInstanceId",
    "workflowRunningId",
    "recordType",
    "taskId",
    "taskKind",
    "attributeId",
    "automationFlowId",
    "contactId",
    "campaignActivityId",
    "sendId"
  ]
}
k
Can you go to the cluster manager UI > Table > Segments and check where if this segment is present in
server-0
or not -
workflowEvents__1__0__20220720T1533Z
a
It is present
let me know if u need any more details @Kartik Khare We are using upsert so there are some duplicates with eventTimestamp as well
k
Yeah, I am looking at the code if there is some edge case or not. The response with 2 rows is the correct one right?
a
yes the 2 rows is the correct one ideally i wanted to softdelete the second row but becuase the qery somestimes does not return it it is not gettign deleted
k
Can you also check what gets returned from server-1
a
Server 1 is only 1 node
Copy code
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	softDelete
cdp-dl-pinot-k8s-server-1.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2023-01-19 21:11:54.709	1642972105000	workflowEvents__1__191__20230119T1352Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:Shopify:1003	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	a9a0c958-39da-4d39-932e-118214590009	1003	null
k
Interesting, so both
server-0
and
server-1
have incorrect state somehow. Can you send me the result of the following query from all 3 servers
Copy code
select $hostName, ToDateTime(eventTimestamp, 'yyyy-MM-dd HH:mm:ss.SSS') as eventTimestamp,
createdOn, 
$segmentName, 
"accountId", 
"workflowDefinitionId", 
"workflowInstanceId", 
"workflowRunningId", 
"recordType", 
"taskId", "taskKind", "attributeId", "automationFlowId", "contactId", "campaignActivityId", "sendId",orderNumber, softDelete  
from workflowEvents 
where accountId = 1100678423876 
AND campaignActivityId = '528680c5-6b72-4e7d-8a9b-febe061d28e7'
AND workflowDefinitionId = 'c746f8df-c25a-4768-9093-8a0f96596022'
AND workflowInstanceId = '259e60b1-3e21-4db3-95a1-33b98ae02e4b'
AND workflowRunningId =  '0e811aff-1dfe-411e-b942-6ba495b6a661'
AND recordType = 'attribution' 
AND taskId = '12ed81dc-0830-43d1-9c0c-50094f1d6604'
AND taskKind = 'grpc'
AND attributeId = 'gif_order:cbbe27299b349eb5ff1c71d9d7ec3416'
AND automationFlowId = 'a66ebfc7-213a-474d-a606-b318646dcd29'
AND contactId = '6b27480e-7c89-11ec-b2bc-fa163ef30863'
AND campaignActivityId = '528680c5-6b72-4e7d-8a9b-febe061d28e7'
AND sendId = 'null'
Might have to replace sendId = 'null' with
sendId IS NULL
a
Copy code
cdp-dl-pinot-k8s-server-0
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	activityType	softDelete
cdp-dl-pinot-k8s-server-0.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__55__20220913T1548Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	true
Copy code
cdp-dl-pinot-k8s-server-1
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	activityType	softDelete
cdp-dl-pinot-k8s-server-1.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__55__20220913T1548Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	true
Copy code
cdp-dl-pinot-k8s-server-2
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	activityType	softDelete
cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__0__20220720T1533Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	null
k
So the row for this particular key def exists and with exact same
createdOn
timestamps although in different segment. What is the activityType in all of these, I missed it in the query
a
server-2 is different
I updated the above results with activityType
yes the data is same
k
Is the last column softDelete or acitivityType? I thought it was all
null
before
a
last column is softDelete
k
Hmm, it is weird a row with exact same Primary key and same eventTimestamp and createdOn has
softDelete
null in server-2 while
true
in
server-1
and
server-0
Can you run the same query with
option(skipUpsert = true)
I am suspecting some events got missed by server-2
a
Copy code
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	activityType	softDelete
cdp-dl-pinot-k8s-server-1.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__0__20220720T1533Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	null
cdp-dl-pinot-k8s-server-1.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__55__20220913T1548Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	true
Copy code
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	activityType	softDelete
cdp-dl-pinot-k8s-server-0.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__0__20220720T1533Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	null
cdp-dl-pinot-k8s-server-0.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__55__20220913T1548Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	true
`
Copy code
$hostName	eventTimestamp	createdOn	$segmentName	accountId	workflowDefinitionId	workflowInstanceId	workflowRunningId	recordType	taskId	taskKind	attributeId	automationFlowId	contactId	campaignActivityId	sendId	orderNumber	activityType	softDelete
cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__0__20220720T1533Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	null
cdp-dl-pinot-k8s-server-2.cdp-dl-pinot-k8s-server-headless.dp-metrics-pinot.svc.cluster.local	2022-01-23 21:21:10.786	1642972105000	workflowEvents__1__55__20220913T1548Z	1100678423876	c746f8df-c25a-4768-9093-8a0f96596022	259e60b1-3e21-4db3-95a1-33b98ae02e4b	0e811aff-1dfe-411e-b942-6ba495b6a661	attribution	12ed81dc-0830-43d1-9c0c-50094f1d6604	grpc	gif_order:cbbe27299b349eb5ff1c71d9d7ec3416	a66ebfc7-213a-474d-a606-b318646dcd29	6b27480e-7c89-11ec-b2bc-fa163ef30863	528680c5-6b72-4e7d-8a9b-febe061d28e7	null	null	null	true
Looks like the same
the order of events might be different in server2 vs server1 and server0 as both events have the same eventTimestamp so the latest in order wins
k
That is unlikely because the record is present in first and 56th segment in all servers. Even if the comparisonColumn values is the same (
eventTimestamp
here), we use the segment name for the comparison
But something happened along those lines only, but reloading should have solved that issue then.
a
i see yeah good point the segment is different ..i am using 0.9.1 pinot..was there any fix made on the later versions in reload ?
m
Yes I recommend upgrading to latest, as a few fixes were made in this area
k
Before upgrading, you can simply try restarting the server-2.
a
i did restart all 3 servers yesterday
I was thinking of upgrading when 0.12.0 is released as that has a fix for GetKinesis Throttling .What is the ETA for that release ?
k
It should be released in 1 week I think. The rc0 has already been cut in the repo.
a
ok thanks
r
@Abhijeet Kushe iiuc, you're saying that the record with taskId
gif_order:cbbe27299b349eb5ff1c71d9d7ec3416
is missing from server-0 and server-1. As you have
"shardIteratorType": "LATEST"
in your tableConfig, could it be the case that server-2 consumed this older record (
eventTimestamp=*2022*-01-23 21:21:10.786 1642972105000
) before server-0 and server-1 began consuming from this topic? Did you perhaps increase your
replicasPerPartition
recently?
a
Both records are present @robert zych in all 3 nodes but 1 node displays it differently And to answer your question we did increase replicas
r
Interesting. Your
skipUpsert=true
results for server-0 and server-1 do include
taskId=gif_order:cbbe27299b349eb5ff1c71d9d7ec3416
but but with a duplicated
eventTimestamp
and different values for
softDelete
(ie
null
,
true
). Perhaps server-0 and server-1 processed these records in a different order than server-2? I noticed that your schema includes
eventTimestampNanos
, assuming that it resolves the duplication, could use that as your
comparisonColumn
in your
upsertConfig
and reload your segments?
a
Yes we added event timestampNanos but unfortunately cannot modify our table to use that field so it is only for triaging Yes I think ordering is probably the reason but I am confused I thought the segment replication should ensure the same segment exists on each server so why does the order matter
r
Not sure why order matters (non-deterministic sort?). If you can't change the current table, can you create a new realtime table and re-process the older events (assuming they're still in Kinesis)?
a
Order does matter as we soft delete events and we cannot create a new table
r
If you can't create a new realtime table because Kinesis no longer has the old events, then can you create an offline table and use the hybrid table approach instead?
a
No we are using upsert and only realtime supports it hybrid does not...but why are you suggesting creating a new table can I ask ?.
r
With the hybrid table approach you could setup an ETL job that would dedup your historical data and upload to an offline table with the same name.
a
No hybrid is not supported for Upsert
r
it's not. just an idea as an alternative to upgrading if you run into issues during/after the upgrade. one the problems with this workaround is potentially having duplicates if your query crosses over the time boundary and the same pks exists in both the realtime and offline tables
a
No it is not an alternative for the exact problem you suggested.We have discussed with the Pinot team in the past and have ruled out that option.If the problem is not solved after the upgrade then we have to find out the root cause.The right solution is to have distinct with pagination then we can solve this problem but that will take time
👍 1
I deleted the server-2 pod and once it came back I found that Server-0 now started showing the 2 values whereas server-1 and server-2 are now showing only 1 value
🤯 1
r
With upsert on or off?
a
Upsert on
r
I'm testing version 0.11.0 with duplicate timestamps and multiple replicas and finding that with upsert turned on the query results are consistent and contain the version with the greatest offset
a
Can you downgrade to 0.9.1 with the same data and see if that consistency is lost ?
r
Will do
Version 0.9.1 appears to be behaving the same as version 0.11.0
In my testing I didn't change
replicasPerPartition
and I have
"stream.kafka.consumer.prop.auto.offset.reset": "smallest"
a
I am using kinesis is that option available for kinesis is is it just Kafka ?
r
I'm using Kafka
a
Ok
r
From an upsert perspective, it doesn't matter if Kafka or Kinesis is being used
I might be able to reproduce the consistency issue if I re-test with
"stream.kafka.consumer.prop.auto.offset.reset": "largest"
and increasing
replicasPerPartition
a
is that value applicable for kinesis
r
stream.kafka.consumer.prop.auto.offset.reset
should be the equivalent of your
shardIteratorType
for Kinesis. Now I'm seeing a consistency issue for the updated record. Here are the steps I performed: 1. Created the kafka topic 2. Created the table with replicasPerPartition=1 3. Started Server-1 4. Produced records to the topic 5. Ran query and got all the records 6. Stopped Server-1 7. Started Server-2 8. Increased replicasPerPartition to 2 9. Ran query and got no records 10. Published 1 record to the topic (update to existing pk) 11. Ran query and got 1 record (latest version of updated record) 12. Stopped Server-2 13. Started Server-1 14. Ran query and got all the original records (update wasn't processed)
Should I proceed to run these steps against version 0.11.0?
a
i did increase my replicas from 1 to 3 when i added nodes ..so yes this is closer what happened my case
yes please try these steps with 0.11.0
r
version 0.11.0 behaves the same as version 0.9.1
a
ok thanks
r
Your welcome. I’ll also run the test against master
a
I think it makes sense to modify the timestamp by 1 for soft deletes to ensure it is not order dependent as we don't know which replica is written in what fashion..
But I do have question as to why the duplicates switched from server2 to server0 ..it is almost as the primary node handles this differently from replicas meaning there is a flag to indicate replica v/s primary. ..if primary server is hit then it is a different order and replica is a different order
m
No, there is no primary / replica concept for servers. All servers consume independently
a
I see so if there are 2 shards in kinesis and we have 3 nodes do all 3 consume from both shards ?
m
Shards will be distributed across servers and replicas (no master / slave).
So if you had 4 shards replication of 2 and 4 servers then one server will consume 2 shard each
a
Ok we have 2 shards and 3 servers with replication 3 so 3 nodes will consume both
m
Yes
r
if you have multiple shards, then for upsert to work correctly, a primary key must be in a single shard and not present in both shards
m
Yes I believe that is already the case for their setup.
a
Yes it is
r
Increasing the timestamp by 1 for a soft delete is a good idea because if events are produced out-of-order (eg soft delete before another update) then upsert will ignore the other update and prevent accidental restores
a
Yeah I am thinking of doing that ... regarding the distinct support for select with pagination...any ideas when that will be released?
r
Server-2 isn't pickup any messages from the topic now. I'm running the latest version from the master branch via IntelliJ. If I can't get that resolved then I'll have to wait for the 0.12.0 release.
👍 1
I forgot to mention that Server-1 did process the update after re-start. If you recall, this wasn't what I observed in earlier versions (0.9.1, 0.11.0).
a
Which version are you using now?
r
Latest from the master branch which I think will be 0.12.0 soon
a
ok and are u seeing the duplicate data issue ?
r
No, I didn’t find a duplicate. After server-1 was restarted it processed the update correctly.
Server-2 is working now. Perhaps it was a timing issue. Not seeing any consistency issue with upsert in the latest version.
👍 1
a
That is great would be nice to try the upgrade this week
👍 1
r
The only inconsistencies I'm seeing now is that Server-2 doesn't have the older records
But this is to be expected under this configuration
To reduce these types of inconsistencies I would recommend changing your
shardIteratorType
from LATEST to TRIM_HORIZON
a
There is this change as well which was done in 2021
this might impact consuming segments but not segments that are closed Both the records I had were from closed segments
r
Interesting. Which shardIteratorType does it default to if the sequenceNumber is null?
It would have to default to either LATEST or TRIM_HORIZON. Either way you might see inconsistencies between servers. It depends on your Kinesis and Pinot retention settings
This is always the case when your adding replicas and your Pinot retention settings is greater than your Kinesis/Kafka retention settings. Does that make sense?
a
Default has to what I have set.I thought Pinot maintains the sequence number in zookeeper.Where are retention settings configured ? So yes the newly added replica will pick up latest whereas the other should be sequence if it is implemented this way.But if that is the case that data would be missing not out of order
👍 1
r
Your Pinot retention setting are
"retentionTimeUnit": "DAYS",
and
"retentionTimeValue": "1826",
. I'm not familiar with Kinesis retention settings.
a
yeah retention is definitely more as we dont use hybrid.But this is a good catch.May be @Kartik Khare can let us know if the new replicas will pick from the sequence number or from LATEST
👍 1
@Kartik Khare can you please take a look?
k
the behavior is independent of the plugin. Whatever happens for Kafka will be the same for Kinesis.
a
Ok so I guess there is a potential for data inconsistency based on what @robert zych mentioned the other replicas will read from LATEST and not sequence number
r
the exact behavior (smallest or largest) of the newer replicas isn't that important if the retention of Kafka/Kinesis is much shorter. this inconsistency is visible only when the newer replicas are online and your query requires older data.
the question I have is how to get very old data on the new replicas?
a
Max Retention of kinesis is 7 days in AWS so Pinot retention to be smaller than that is highly unlikely
r
now I remember. when changing
replicasPerPartition
you have to run a rebalance
a
Yes I did run rebalance