This message was deleted.
# troubleshooting
s
This message was deleted.
s
Could you share your ingestion spec and the failed task log?
j
My ingestion spec is here:
Copy code
{
  "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
    }
  }
}
Unfortunately we have had trouble getting the task logs before. When we go to the taks logs we only get garbage heap info like:
Copy code
[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] ...
s
Not sure about why the task log is giving you that. But a few items regarding your ingestion spec: • I don't see a
taskDuration
, 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?
j
Thanks @Sergio Ferragut I did try
lateMessageRejectionPeriod
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?
It seems like the first job that fails, fails with overlord logs like:
Copy code
2022-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]
My best explanation of this log is that 1. A late arriving task gets picked up. 2. The late arriving task can not get properly allocated 3. The task not getting properly allocated kills the running task 4. Next tasks get failed as well
However, I could be off base here
s
The problem with late data is that it will create another segment for the corresponding time interval, which will likely already have been submitted to deep storage which just means you end up creating many small segments. In order to avoid that, you setup a threshold
lateMessageRejectionPeriod
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.
j
Great, thanks for the advice 👍
g
these "Could not allocate segment for row" errors can occur if you have mixed granularities in the same datasource… does this seem likely in your situation?
we've recently adjusted the logic here to be more robust, as well; you can try updating to 24.0 if you are not already on it. if you are on it then i hope we figure out what is going on in your situation, since i'd like to fix whatever it is!
j
@Gian Merlino we have a compaction task which I think might changes granularity, but that would be after 10 days. The segment here on 11/11 is only 3 days older than the one thats trying to run on 11/4 and should be using the same granularity. In any case though, adding
lateMessageRejectionPeriod
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.
s
@James Von Kaenel since the
lateMessageRejectionPeriod
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?
j
@Sergio Ferragut I don’t think that is the case, because we had
lateMessageRejectionStartDateTime
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:
Copy code
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]
]
s
What granularity are the compacted segments on?
j
They have a daily segment granularity