does anyone have experience setting up the <offlin...
# troubleshooting
l
does anyone have experience setting up the offline pinot managed flows? I have 2 questions… 1. is this the same as moving completed realtime segments to offline and is this config required for it to work
Copy code
},
  "tenants": {
    "broker": "DefaultTenant",
    "server": "DefaultTenant",
    "tagOverrideConfig": {
      "realtimeCompleted": "DefaultTenant_OFFLINE"
    }
  }
2. is there anywhere I can see a log that this is in fact working I have setup the configs but I’m unsure as to how to tell it’s doing what it’s supposed to be doing 3. documentation is a little bit misleading in the sense of the new updates that we have done to pinot as well as different examples doing different things that are not explained in the documentation
m
Yes, this will move segments from realtime to offline.
@Neha Pawar could you help with the other questions?
n
the config you’d posted above is not the same thing. this tagOverrideConfig only moves only different servers, but still keeps everything in realtime table
🙌 1
you should see logs in controller, every time the task is scheduled. And in the minion logs, everytime the task is executed
l
hmmm
for some reason i don’t see the minion logging at all
Copy code
{
  "REALTIME": {
    "tableName": "ads_metrics_dev_REALTIME",
    "tableType": "REALTIME",
    "segmentsConfig": {
      "schemaName": "ads_metrics_dev",
      "retentionTimeUnit": "DAYS",
      "retentionTimeValue": "7",
      "replication": "1",
      "timeColumnName": "serve_time",
      "allowNullTimeValue": false,
      "replicasPerPartition": "1"
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant",
      "tagOverrideConfig": {
        "realtimeCompleted": "DefaultTenant_OFFLINE"
      }
    },
    "tableIndexConfig": {
      "invertedIndexColumns": [],
      "noDictionaryColumns": [
        "click_count",
        "order_count",
        "impression_count",
        "cost",
        "revenue"
      ],
      "streamConfigs": {
        "streamType": "kafka",
        "stream.kafka.topic.name": "ads-stats",
        "stream.kafka.broker.list": "<http://kafka.dev.com:9092|kafka.dev.com:9092>",
        "stream.kafka.consumer.type": "lowlevel",
        "stream.kafka.consumer.prop.auto.offset.reset": "largest",
        "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": "250M"
      },
      "enableDynamicStarTreeCreation": false,
      "aggregateMetrics": true,
      "nullHandlingEnabled": false,
      "rangeIndexColumns": [],
      "rangeIndexVersion": 2,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "sortedColumn": [
        "shop_id"
      ],
      "bloomFilterColumns": [
        "shop_id",
        "listing_id"
      ],
      "loadMode": "MMAP",
      "onHeapDictionaryColumns": [],
      "varLengthDictionaryColumns": [],
      "enableDefaultStarTree": false
    },
    "metadata": {},
    "quota": {},
    "task": {
      "taskTypeConfigsMap": {
        "RealtimeToOfflineSegmentsTask": {
          "bucketTimePeriod": "1d",
          "bufferTimePeriod": "2d",
          "roundBucketTimePeriod": "1h",
          "mergeType": "concat",
          "maxNumRecordsPerSegment": "5000000"
        }
      }
    },
    "routing": {},
    "query": {},
    "fieldConfigList": [],
    "ingestionConfig": {},
    "isDimTable": false
  }
}
this is my realtime table config
and this is my offline table config
Copy code
{
  "OFFLINE": {
    "tableName": "ads_metrics_dev_OFFLINE",
    "tableType": "OFFLINE",
    "segmentsConfig": {
      "schemaName": "ads_metrics_dev",
      "retentionTimeUnit": "DAYS",
      "retentionTimeValue": "186",
      "replication": "1",
      "segmentPushType": "APPEND",
      "timeColumnName": "serve_time",
      "allowNullTimeValue": false,
      "replicasPerPartition": "1",
      "segmentPushFrequency": "HOURLY"
    },
    "tenants": {
      "broker": "DefaultTenant",
      "server": "DefaultTenant"
    },
    "tableIndexConfig": {
      "invertedIndexColumns": [],
      "noDictionaryColumns": [
        "click_count",
        "order_count",
        "impression_count",
        "cost",
        "revenue"
      ],
      "enableDynamicStarTreeCreation": false,
      "aggregateMetrics": false,
      "nullHandlingEnabled": false,
      "rangeIndexColumns": [],
      "rangeIndexVersion": 2,
      "autoGeneratedInvertedIndex": false,
      "createInvertedIndexDuringSegmentGeneration": false,
      "sortedColumn": [
        "shop_id"
      ],
      "bloomFilterColumns": [
        "shop_id",
        "listing_id"
      ],
      "loadMode": "MMAP",
      "onHeapDictionaryColumns": [],
      "varLengthDictionaryColumns": [],
      "enableDefaultStarTree": false
    },
    "metadata": {},
    "quota": {},
    "routing": {},
    "query": {},
    "fieldConfigList": [],
    "ingestionConfig": {},
    "isDimTable": false
  }
}
I added
controller.task.frequencyInSeconds=3600
to my configmap.yaml in the controller
any way i can check if those changes were picked up in the logs
bumping this looking for some help on this regard
m
Did you directly edit
controller.task.frequencyInSeconds=3600
on the pod, or you restarted the pod? For controller to pick it up, it has to be restarted.
l
i edited, pushed and then restarted the controller
m
I see the following in the code:
Copy code
// Deprecated as of 0.8.0
    @Deprecated
    public static final String DEPRECATED_TASK_MANAGER_FREQUENCY_IN_SECONDS = "controller.task.frequencyInSeconds";
    public static final String TASK_MANAGER_FREQUENCY_PERIOD = "controller.task.frequencyPeriod";
@Neha Pawar which one of these is the latest, and is the doc updated with this?
l
is there any endpoint that can help see if they are at least schedule? i was trying out some of then but no luck
n
checking
Do you see the following logs in controller:
Copy code
<http://LOGGER.info|LOGGER.info>("Adding periodic task: {}", periodicTask);
and
Copy code
<http://LOGGER.info|LOGGER.info>("Starting periodic task scheduler with tasks: {}", _tasksWithValidInterval);
and
Copy code
<http://LOGGER.info|LOGGER.info>("Starting {} with running frequency of {} seconds.", periodicTask.getTaskName(),
and if yes do you see PinotTaskManager listed in it?
here’s one thing I would try. There’s unfortunately 2 ways of enabling minion tasks. One is with the
controller.task.frequencyInSeconds
config, which you have tried. Other is, adding a
schedule
directly into the task config. So, do these 2 steps 1. add
schedule
to you config, like this
Copy code
"RealtimeToOfflineSegmentsTask": {
          "bucketTimePeriod": "<your value>",
          "bufferTimePeriod": "<your value>",
          "schedule": "0 0/10 0 ? * * *"
        }
2. add this config to controller cofig
Copy code
controller.task.scheduler.enabled=true
l
thank you gonna check this out
yep I do see it
Copy code
Starting periodic task scheduler with tasks: [Task: PinotTaskManager, Interval: 3600
trying the other things you said
hmm i do see this in one of our controllers
Copy code
Start generating task configs for table: ds_metrics_dev_REALTIME for task: Realtim
eToOfflineSegmentsTask
Window with start: 5919648480000000 and end: 5919648501600000 is not older than buffer 
time: 25200000 configured as 7h ago. Skipping task generation: RealtimeToOfflineSegment
sTask
n
what is this format of time column?
5919648480000000
l
Copy code
"dateTimeFieldSpecs": [
    {
      "name": "serve_time",
      "dataType": "LONG",
      "format": "1:HOURS:EPOCH",
      "granularity": "1:HOURS"
    }
  ]
does that help?
n
yes, the time column looks like it has incorrect values
l
this is the config
Copy code
"task": {
      "taskTypeConfigsMap": {
        "RealtimeToOfflineSegmentsTask": {
          "bucketTimePeriod": "6h",
          "bufferTimePeriod": "7h",
          "mergeType": "concat",
          "maxNumRecordsPerSegment": "5000000",
          "schedule": "0 0/10 0 ? * * *"
        }
      }
    },
how do i make it have the correct values?
n
i mean, unrelated to this config, the time column has incorrect values.
for instance, today’s hoursSinceEpoch value is 456790
l
ohh so that’s misconfigured then
ok so serve_time is the epoch timestamp to the hour
n
but it seems data has values of the range
5919648480000000
. You ca confirm what the range is by doing
select min(time column), max(time column) from table
l
so data form serve_time looks liek this
1644462000
n
this is seconds since epoch ^^
l
Copy code
1644433200
n
is this a production setup or test evironment?
l
we are testing
so should it be seconds then?
n
can you recreate table with proper schema cofnig then? this misconfiguration must’ve affected all completed segment’s zk metadata
l
yea we can
n
Copy code
"dateTimeFieldSpecs": [
    {
      "name": "serve_time",
      "dataType": "LONG",
      "format": "1:SECONDS:EPOCH",
      "granularity": "1:HOURS"
    }
l
so our serve_time i guess seconds but we truncate always to the hour
still seconds but you hourly i gues
does that make sense?
what other things can that mess up just for my own education
n
yup makes sense.
retention will not work. so whatever value you’ve set in retention for the table, your segments will never get deleted, because this calculates to some hours since epoch in the future
l
omg that would be sad weelll thank god it’s dev
what would you do if you have a bad config like that in prod? can you recover?
n
as a thought experiment, i could probably think of some way to make it work. But practically very hard, because those incorrect values are set into the segment. We could have added a new derived column. And done zk edits to force change the time column (which normally is a disallowed operation)
l
ohh so that’s why i see that weird epoch in the segment info
this right?
Copy code
"segment.end.time": "5919959520000000",
n
yup!
l
anything that would impact lookups to the table having a config like this with the wrong format
or segment creation
n
would not have impacted the query. segment creation also is fine (because the value put in was some valid futuristic hoursSinceEpoch). If it was other way around (config has MILLIS but provided hours) it would’ve failed to create segment because it wouldn’t have bee a valid millisEpoch value. OR if the time column was in a simple date format, then it would’ve definitely failed.
i’m trying to think how we could’ve flagged this earlier.
l
i’m issuing fixes today thank you so much Neha, thankfully we haven’t gotten to prod but our shadow data is all corrupted i guess another way would be to recreate the table and somehow fix the data in the segments on the deep store and then recreate them in the new table? is that even possible or get the data from somewhere else (we also synch to bigquery) and fix it and offline push it to an offline table in pinot and once that backfill is done switch tables or something like that
hey Neha does that epoch sense to you?
Copy code
Trying to schedule task type: RealtimeToOfflineSegmentsTask, isLeader: true
Start generating task configs for table: metrics_dev_REALTIME for task: RealtimeToOfflineSegmentsTask
Window with start: 1644451200000 and end: 1644537600000 is not older than buffer time: 172800000 configured as 2d ago. Skipping task generation: RealtimeToOfflineSegmentsTask
does this*
n
yup. looks right now
l
so let me see if i understand this correctly, (btw this is what we setup that we had without the
"schedule"
attribute on the table). the periodic task scheduler is setup to run every 3600 seconds, so every hour so every hour this scheduler wakes up, and figures who has to run, when a controller is restarted this starts by default yes? so based on the message above ^ it will keep on trying to run till the buffer time > than the window start and beginning then it will run. so with this current setup, it means that new data gets to the offline table every 2 days from the day before?
n
Copy code
the periodic task scheduler is setup to run every 3600 seconds, so every hour so every hour this scheduler wakes up, and figures who has to run, when a controller is restarted this starts by default yes? so based on the message above ^ it will keep on trying to run till the buffer time > than the window start and beginning then it will run.
correct
so with this current setup, it means that new data gets to the offline table every 2 days from the day before?
- not sure i follow this sentence. Your window looks like 1 day. And buffer 2 day. So every day, a new window will move to offline. This will start happening after 2 days
l
i guess what i meant to say is say i started ingesting data wednesday, it will start uploading data from wednesday, today?
and then tomorrow it would do thursday?
Window with start: 1644451200000 and end: 1644537600000 is not older than buffer time: 90000000 configured as 25h ago. Skipping task generation: RealtimeToOfflineSegmentsTask
I changed the bugger to 25h given that the latest the data can reach is 1h does this mean that this will run after a day and an hour? cause I don’t see the difference
90000000
increasing as the hours go by
n
Copy code
i guess what i meant to say is say i started ingesting data wednesday, it will start uploading data from wednesday, today? -yes this is right
9000000 won't keep increasing, but the current time will keep increasing. Making difference between current and window end time greater than 9000000 at some point
l
hmm for some reason this still is not working 😞
I should see my offline table grow in size yes?
n
What do the logs say on the controller
l
Copy code
Workflow TaskQueue_RealtimeToOfflineSegmentsTask or job TaskQueue_RealtimeToO
fflineSegmentsTask_Task_RealtimeToOfflineSegmentsTask_1644853530578 is alread
y failed or completed, workflow state (IN_PROGRESS), job state (COMPLETED), c
lean up job IS.
Job: TaskQueue_RealtimeToOfflineSegmentsTask_Task_RealtimeToOfflineSegmentsTa
sk_1644853530578 has either finished already, never been scheduled, or been r
emoved from DAG
i do see this as well
Copy code
Minion_pinot-minion-0.pinot-minion-headless.pinot.svc.cluster.local_9514 tra
nsit TaskQueue_RealtimeToOfflineSegmentsTask_Task_RealtimeToOfflineSegmentsTa
sk_1644853530578.TaskQueue_RealtimeToOfflineSegmentsTask_Task_RealtimeToOffli
neSegmentsTask_1644853530578_0|[] from:RUNNING to:TASK_ERROR, relayMessages: 
0
ooooo
Copy code
Caught exception while fetching segment from: <gs://pinot-data/me>
trics_dev/metrics_dev__0__0__20220210T1525Z to: /var/pinot/minion/dat
a/RealtimeToOfflineSegmentsTask/tmp-06a47855-b7e0-4dd2-9682-1f5bb2108d3f/tarr
edSegmentFile_0
java.lang.IllegalStateException: PinotFS for scheme: gs has not been initiali
zed
is it cause i have to put some options about the deep store on the minion cause the minion interacts w/ the deep store
l
ok one i fix that what should I expect
once*
the behavior to be
will try to catch up to present day?
omg it’s working 😭
😄 1
now i have more questions tho 😄
so i see data being moved from last thursday my config is the following:
Copy code
"taskTypeConfigsMap": {
        "RealtimeToOfflineSegmentsTask": {
          "bucketTimePeriod": "24h",
          "bufferTimePeriod": "25h",
          "roundBucketTimePeriod": "1h",
          "mergeType": "concat",
          "maxNumRecordsPerSegment": "5000000"
        }
      }
i see this segment on my offline table
metrics_dev_1644505200_1644534000_0
so that’s 3pm till 11pm gmt
1644505200_1644534000
based on these timestamps
what does that mean
i was expecting to see something more around ok this is from beginning of day till end of date timestamps
n
it could be that your realtime table just has data from Feb 10th 3pm to Feb 10th 11pm. The job will always look for the entire window [Feb 10th 00 to Feb 11th 00)
the log statement which tells you the windoe it picked, should always match the day boundary
l
yea that makes sense this is our dev env so i wouldn’t except anyone touching stuff after certain point
now, this is the data from thursday, what ways do i have to actually make it run for other days?
like i want data from friday also to be moved
since it’s now monday
does that make sense
n
It should happen automatically, once the new window is past the buffer
You can keep looking at that log in the controller which printed window and buffer
l
right
is it normal that the older segments are still listed in the realtime table?
like in the UI I see this
metrics_dev__0__0__20220210T1525Z
in the realtime table
and that segment should now be in the offline table
also one of the segments looks super weird in the offline table =
metrics_dev_1644624000_1644624000_0
that’s the same timestamp o.o and it has 8 docs
n
try a query like
select min(<time column>), max(<time column>), $segmentName from <table>_OFFLINE group by $segmentName order by $segmentName
It should show you that the min and max indeed are what you see in the segmentName
👀 1
segments aren’t actually moving from realtime to offline. The data from realtime is being used to fill up the offline table. This becomes a hybrid table. Segments can exist in both. They will be deleted from realtime table according to the retention set in the realtime table. Queries should be fine, because Pinot uses a time boundary https://docs.pinot.apache.org/basics/components/broker
l
ohh so they will exist on both till it hits retention but when actually querying for them the broker will know not to fetch it from the realtime table but the offline one based on the time boundary.
another question that I have does the
segmentPushFrequency
has to match in this case how often we push this offline segments?
for the query that you showed above, for that particular segment i see the min and max as the same
1644624000
there’s a segment that was collecting data from that friday that did from 1pm to 9pm gmt ohhh i see cause then some data was processed on saturday so we got some data from saturday
1644624000
but is all GMT, so i guess that makes sense, i just got confused cause i’m on EST
so i think it’s cool
timezones 🥴
n
you can ignore
segmentPushFrequency
for your case. that’s only relevant for manual offline segment uploads
l
thank you 🙏
so now that i see it working in dev i’m gonna try to apply this change to our shadow prod dataset
hey @Neha Pawar we are having this working with staging data now, one question that i had, during the weekend the controller tried to move segments, but it failed cause the minion machine was running out of disk, so today i incremented the size disk of that pod, so it moved the segments that happened on friday, am i correct to assume that in the next hour it will move the ones on saturday and keep going till it doesn’t have more to run?
n
yes, you are correct