Slackbot
11/14/2022, 7:18 PMSergio Ferragut
11/14/2022, 9:11 PMJames Von Kaenel
11/14/2022, 9:15 PM{
"type": "kinesis",
"spec": {
"dataSchema": {
"dataSource": "data_source",
"metricsSpec": [
{
"name": "ex_count",
"type": "count"
}
],
"granularitySpec": {
"segmentGranularity": "hour",
"queryGranularity": "minute",
"rollup": true,
"type": "uniform"
},
"dimensionsSpec": {
"dimensions": [
"ex.dimension"
},
"timestampSpec": {
"column": "at",
"format": "posix"
}
},
"ioConfig": {
"stream": "<kinesis-stream-name>",
"deaggregate": true,
"useEarliestSequenceNumber": true,
"inputFormat": {
"type": "json",
"flattenSpec": {
"useFieldDiscovery": true,
"fields": [
{
"expr": "$.message.ex_field",
"name": "ex_field",
"type": "path"
}
]
}
},
"endpoint": "kinesis.<kinesis-region>.<http://amazonaws.com|amazonaws.com>",
"taskCount":"<taskCount>",
"replicas":"<replicas>",
"lateMessageRejectionStartDateTime": "2022-11-10T08:15Z",
"type": "kinesis"
},
"tuningConfig": {
"resetOffsetAutomatically": false,
"skipSequenceNumberAvailabilityCheck":true,
"type": "kinesis",
"reportParseExceptions": false,
"logParseExceptions": true,
"intermediatePersistPeriod": "PT10M",
"maxRowsInMemory": 75000
}
}
}James Von Kaenel
11/14/2022, 9:16 PM[GC Worker End (ms): Min: 9088.7, Avg: 9088.7, Max: 9088.7, Diff: 0.0]
[Code Root Fixup: 0.0 ms]
[Code Root Purge: 0.0 ms]
[Clear CT: 0.1 ms]
[Other: 0.4 ms]
[Choose CSet: 0.0 ms]
[Ref Proc: 0.1 ms] ...Sergio Ferragut
11/14/2022, 9:50 PMtaskDuration, probably a good idea to set this to 1 hour, or maybe 2. This means that each task will stop after that time to output segments and another task will be started to continue reading from the same shards.
• I have not seen lateMessageRejectionStartDateTime before, usually it is done with lateMessageRejectionPeriod which you set to a duration like PT1H such that messages arriving with a timestamp that is more than an hour before the task is created are rejected. Is this what you mean to do?James Von Kaenel
11/14/2022, 9:52 PMlateMessageRejectionPeriod and it worked but I was worried that it would cut off some data, and thought that lateMessageRejectionStartDateTime would cut off data intially, but not afterwards.
However, i’m guessing that we probably should expect much data to arrive over full hour late, so lateMessageRejectionPeriod should usually be acceptable?James Von Kaenel
11/14/2022, 10:01 PM2022-11-14T08:06:30,766 INFO [org.apache.druid.metadata.SqlSegmentsMetadataManager-Exec--0] org.apache.druid.metadata.SqlSegmentsMetadataManager - Polled and found 28,486 segments in the database
2022-11-14T08:06:36,375 INFO [qtp1329043305-79] org.apache.druid.indexing.overlord.TaskLockbox - Added task[index_kinesis_data_source_6a6dc7bae0c9a5b_fkbelhbe] to TaskLock[TimeChunkLock{type=EXCLUSIVE, groupId='index_kinesis_data_source', dataSource='data_source', interval=2022-11-11T21:00:00.000Z/2022-11-11T22:00:00.000Z, version='2022-11-14T08:06:36.263Z', priority=75, revoked=false}]
2022-11-14T08:06:36,375 INFO [qtp1329043305-79] org.apache.druid.indexing.overlord.MetadataTaskStorage - Adding lock on interval[2022-11-11T21:00:00.000Z/2022-11-11T22:00:00.000Z] version[2022-11-14T08:06:36.263Z] for task: index_kinesis_data_source_6a6dc7bae0c9a5b_fkbelhbe
2022-11-14T08:06:36,386 ERROR [qtp1329043305-79] org.apache.druid.indexing.common.actions.SegmentAllocateAction - Could not allocate pending segment for rowInterval[2022-11-11T21:49:00.000Z/2022-11-11T21:50:00.000Z], segmentInterval[2022-11-11T21:00:00.000Z/2022-11-11T22:00:00.000Z].
2022-11-14T08:06:37,047 INFO [NodeRoleWatcher[PEON]] org.apache.druid.curator.discovery.CuratorDruidNodeDiscoveryProvider$NodeRoleWatcher - Node[<http://stage-druid-middle-manager-1:8100>] of role[peon] went offline.
2022-11-14T08:06:37,049 INFO [ServerInventoryView-0] org.apache.druid.client.BatchServerInventoryView - Server Disappeared[DruidServerMetadata{name='stage-druid-middle-manager-1:8100', hostAndPort='stage-druid-middle-manager-1:8100', hostAndTlsPort='null', maxSize=0, tier='_default_tier', type=indexer-executor, priority=0}]
2022-11-14T08:06:38,271 INFO [Curator-PathChildrenCache-1] org.apache.druid.indexing.overlord.RemoteTaskRunner - Worker[stage-druid-middle-manager-1:8091] wrote FAILED status for task [index_kinesis_data_source_6a6dc7bae0c9a5b_fkbelhbe] on [TaskLocation{host='stage-druid-middle-manager-1', port=8100, tlsPort=-1}]
2022-11-14T08:06:38,271 INFO [Curator-PathChildrenCache-1] org.apache.druid.indexing.overlord.RemoteTaskRunner - Worker[stage-druid-middle-manager-1:8091] completed task[index_kinesis_data_source_6a6dc7bae0c9a5b_fkbelhbe] with status[FAILED]James Von Kaenel
11/14/2022, 10:08 PMJames Von Kaenel
11/14/2022, 10:09 PMSergio Ferragut
11/14/2022, 10:43 PMlateMessageRejectionPeriod which can be whatever you want, 1 hour or 6 hours or more, it's up to you.
I think the taskDuration is more important so that segments are output to deep storage at least once an hour and a new task starts building new segments at that point. I'm still not sure what the source your failure is though. Each task will request a segment allocation from the overlord when it sees a timestamp that doesn't fall into the intervals for the segments that it is currently building. Long running tasks would tend to have many more segments in progress, so controlling the task duration might help.
It is always a good practice to follow up a streaming ingestion with a compaction or auto-compaction job to consolidate all those potentially small segments that created by the parallel tasks and/or late arriving messages.James Von Kaenel
11/14/2022, 11:49 PMGian Merlino
11/15/2022, 6:32 AMGian Merlino
11/15/2022, 6:33 AMJames Von Kaenel
11/15/2022, 9:21 PMlateMessageRejectionPeriod
seems to make the tasks run smoothly. The data is also looking good.
We are on 0.18.1 so the recent fixes would not apply.Sergio Ferragut
11/15/2022, 9:50 PMlateMessageRejectionPeriod seems to solve it, Is it likely that you were getting timestamps that were older than 10 days which were trying to allocate a segment in one granularity that overlapped with one of the compacted granularity segments?James Von Kaenel
11/15/2022, 9:52 PMlateMessageRejectionStartDateTime applied which was set to only 4 days back from the error (11/10). Also the segment that seems to be giving trouble is from 11/11:
Could not allocate pending segment for rowInterval[2022-11-11T21:49:00.000Z/2022-11-11T21:50:00.000Z], segmentInterval[2022-11-11T21:00:00.000Z/2022-11-11T22:00:00.000Z]
]Sergio Ferragut
11/15/2022, 10:03 PMJames Von Kaenel
11/15/2022, 10:47 PM