Hello folks, I am struggling with hybrid table but...
# troubleshooting
y
Hello folks, I am struggling with hybrid table but i have trouble to make it work. My configs are like below. Offline table data is only displayed without streaming data blended. I could find log of consuming kafka events but couldn’t find broker, controller or server error. My hybrid table creation testing is running on minikube, with pinot 0.9.0. Process I did was creating offline, and then realtime. • bin/pinot-admin.sh AddTable -tableConfigFile hybrid_realtime.json -schemaFile hybrid_schema.json -exec • bin/pinot-admin.sh AddTable -tableConfigFile hybrid_offline.json -schemaFile hybrid_schema.json -exec 1. hybrid_schema.json
Copy code
{
  "schemaName": "transcript",
  "dimensionFieldSpecs": [
    {
      "name": "studentID",
      "dataType": "INT"
    },
    {
      "name": "firstName",
      "dataType": "STRING"
    },
    {
      "name": "lastName",
      "dataType": "STRING"
    },
    {
      "name": "gender",
      "dataType": "STRING"
    },
    {
      "name": "subject",
      "dataType": "STRING"
    },
    {
      "name": "doNotFailPlease",
      "dataType": "STRING",
      "defaultNullValue": ""
    },
    {
      "name": "ts2",
      "dataType": "TIMESTAMP"
    }
  ],
  "metricFieldSpecs": [
    {
      "name": "score",
      "dataType": "FLOAT"
    }
  ],
  "dateTimeFieldSpecs": [
    {
      "name": "ts",
      "dataType": "TIMESTAMP",
      "format": "1:SECONDS:EPOCH",
      "granularity": "1:SECONDS"
    }
  ],
  "primaryKeyColumns": [
    "studentID"
  ]
}
2. hybrid_offline.json
Copy code
{
    "tableName": "transcript_hybrid",
    "tableType": "OFFLINE",
    "segmentsConfig": {
        "replication": 1,
        "timeColumnName": "ts",
        "timeType": "SECONDS"
    },
    "tenants": {
        "broker": "DefaultTenant",
        "server": "DefaultTenant"
    },
    "tableIndexConfig": {
        "loadMode": "MMAP"
    },
    "metadata": {}
3. hybrid_realtime.json
Copy code
{
  "tableName": "transcript_hybrid",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "timeColumnName": "ts",
    "timeType": "SECONDS",
    "schemaName": "transcript",
    "replicasPerPartition": "1"
  },
  "tenants": {},
  "tableIndexConfig": {
    "loadMode": "MMAP",
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.consumer.type": "lowlevel",
      "stream.kafka.topic.name": "transcript",
      "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.stream.kafka.KafkaJSONMessageDecoder",
      "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory",
      "stream.kafka.broker.list": "kafka.local-pinot.svc.cluster.local:9092",
      "realtime.segment.flush.threshold.time": "6h"
    }
  },
  "metadata": {
    "customConfigs": {}
  },
  "routing": {
    "instanceSelectorType": "strictReplicaGroup"
  },
  "upsertConfig": {
    "mode": "FULL"
  }
}
m
do you see any segments are created for the realtime table? You can navigate to the tables/segments via this screen - http://localhost:9000/#/tables
The full logs of each component are at
logs/pinot-all.log
- did you already check for entries in there?
y
it seems like ‘logs/pinot-all.log’ this one got freezed. so i am debugging with stdout from the only server, server-0 a segment is created and stdout show it as consuming but data doesn’t display.
@Mark Needham should i wait until initial flush?
m
I thought you should be able to query while it's consuming, but let's check with @Mayank to be certain. In the meantime you could experiment with the threshold time if it's not too hard to change. I think right now it's set to
6h
- maybe reduce that to a smaller value?
m
Events are ready to be served as soon as they are ingested, segment does not need to be flushed for that.
I am not sure if I follow the original issue, is it that query is returning rows only from offline table?
If logs are missing, then it is log4j setting issue.
If query is returning offline data but not real-time then your time boundary is incorrect (check time value and units). You can check min time in real-time and max time in offline. This could happen if offline has more recent data than real-time (usually due to time unit issue).
y
@Mayank For data return, yes only from offline table
@Mayank I didn’t check the source, but I think its flushing size could make it look like freezed. (red ellipse part)
m
That might just be logging issue. Can you run query on individual tables (add _REALTIME) to table name in the query and see if you get any data
If yes, then most definitely it is time boundary issue
y
@Mayank Following your guide, I tried these cases with but realtime data does not appear. below is type, granularity pair. • TIMESTAMP / ms • TIMESTAMP / s • LONG / ms • LONG /s
m
Can you run select count(*) from transcript_hybrid_REALTIME?
y
Oh.. I got this one
and I can see log from INSTANCES/broker-headless/ERRORS at zk UI
let me summarize
Copy code
"HELIX_ERROR     20211202-132039.000850 STATE_TRANSITION 8c6c680c-f89a-450f-b9c9-a4464d517617": {
      "AdditionalInfo": "Exception while executing a state transition task transcript_hybrid_REALTIMEjava.lang.reflect.InvocationTargetException\n\tat java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)\n\tat java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)\n\tat java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)\n\tat java.base/java.lang.reflect.Method.invoke(Method.java:566)\n\tat org.apache.helix.messaging.handling.HelixStateTransitionHandler.invoke(HelixStateTransitionHandler.java:404)\n\tat org.apache.helix.messaging.handling.HelixStateTransitionHandler.handleMessage(HelixStateTransitionHandler.java:331)\n\tat org.apache.helix.messaging.handling.HelixTask.call(HelixTask.java:97)\n\tat org.apache.helix.messaging.handling.HelixTask.call(HelixTask.java:49)\n\tat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n\tat java.base/java.lang.Thread.run(Thread.java:829)\nCaused by: java.lang.IllegalStateException: Failed to find schema for table: transcript_hybrid_OFFLINE\n\tat shaded.com.google.common.base.Preconditions.checkState(Preconditions.java:518)\n\tat org.apache.pinot.broker.routing.timeboundary.TimeBoundaryManager.<init>(TimeBoundaryManager.java:74)\n\tat org.apache.pinot.broker.routing.RoutingManager.buildRouting(RoutingManager.java:371)\n\tat org.apache.pinot.broker.broker.helix.BrokerResourceOnlineOfflineStateModelFactory$BrokerResourceOnlineOfflineStateModel.onBecomeOnlineFromOffline(BrokerResourceOnlineOfflineStateModelFactory.java:80)\n\t... 12 more\n",
      "Class": "class org.apache.helix.messaging.handling.HelixStateTransitionHandler",
      "MSG_ID": "5173dca5-b009-46cd-b6be-15e57b3a5b01",
      "Message state": "READ"
    }
Copy code
"HELIX_ERROR     20211202-132039.000904 STATE_TRANSITION e7d869d7-0a91-4067-990d-e954ef5228ef": {
      "AdditionalInfo": "Message execution failed. msgId: 5173dca5-b009-46cd-b6be-15e57b3a5b01, errorMsg: java.lang.reflect.InvocationTargetException",
      "Class": "class org.apache.helix.messaging.handling.HelixStateTransitionHandler",
      "MSG_ID": "5173dca5-b009-46cd-b6be-15e57b3a5b01",
      "Message state": "READ"
    },
@Mayank these 2 repeat
m
@Mayank it seems like that line of code assumes that the name of the schema is the same as the table name, which in this case it isn't https://github.com/apache/pinot/blob/master/pinot-broker/src/main/java/org/apache/pinot/broker/routing/timeboundary/TimeBoundaryManager.java#L74
@Yeongju Kang can you try changing the name of your schema to be
transcript_hybrid
to test out that hypothesis?
y
@Mark Needham @Mayank Thank you all! It works!
@Mark Needham By the way, is there any plan to allow different name of table and schema?
m
I think it should allow that - seems like there's a bug. I've created an issue for it - https://github.com/apache/pinot/issues/7856
hope that's ok to include the config you shared?
lemme know if not and I'll update
y
@Mark Needham my pleasure, no worries.
@Mark Needham that was from pinot official doc😁 by the way, what is recommended method to export realtime table's sealed segment as csv, json or avro? I want to extract and upload to s3. After loading, dropping will be my next. Will it be able to be done with swagger api call?
m
I'm not sure how to convert a segment into those formats, but Pinot has the concept of a 'deep store' which keeps compressed copies of segment files in various file formats, of which one is S3. See https://docs.pinot.apache.org/basics/components/deep-store and https://docs.pinot.apache.org/users/tutorials/use-s3-as-deep-store-for-pinot Would that do what you want?
y
@Mark Needham I use ingestion job to import data from s3 like blue part, but couldn’t find red part feature, realtime to s3.
image.png
m
ah ok. So 'segment store' == 'deep store', so the links that I shared should explain how to configure it so that segments are persisted to S3