Hi team, we have some scalability concern of the l...
# ingestion
b
Hi team, we have some scalability concern of the logic used to ingest from SQL-based sources, eg hive. From what I read in the code from
sql_common.py
, it first gets the list of all schemas (L344). Next, for each schema, get the list of tables (L408). Lastly, for each table, get column info (L322). That means, for each run of ingestion, it triggers at least N + M statements against SQL source, e.g. (DESCRIBE <table> in Hive), where N is the number of tables and M is the number of schemas. In our case, we have
over 80K
tables in Hive metastore. Empirically, we tried to ingest one big Hive schema with over 8K tables, and it took 2 hours to finish. And if we scale this duration linearly to 80K tables, that means in our case, each Hive ingestion would take 20 hours to finish, which is not acceptable. What’s your thought or advice on this?
b
run parallel ingest jobs?
b
That could be one way, but there’re overheads. It would mean that we need to create multiple recipes. Each recipe only ingest metadata from a subset of schemas. In case there are new schemas created, we need to add them into one of the existing recipes or create new recipe manually. 🤔
Actually for Hive, one more scalable approach can be getting table & column info by querying Hive Metastore directly instead of doing DESCRIBE statements one by one against Hive. Not sure if anyone from the community takes up this approach before
l
Hi @boundless-student-48844, thanks so much for raising this! I reviewed this thread with the core team earlier today & we agree this is definitely something we need to tackle. I’ve gone ahead and created this feature request; feel free to subscribe to it for updates as we are able to prioritize & begin implementing improvements!
b
Hey Maggie, good to hear and thanks for following up! Definitely look forward to this. The link above is broken. Reattach the link for tracking https://feature-requests.datahubproject.io/b/Developer-Experience/p/optimize-ingestion-for-sql-based-sources-eg-hive
thanks ewe 1
d
Thanks @big-coat-53708, I think we should definitely improve our current Hive ingestion and querying Hive Metastore is better approach and also the parallel processing/partitioned processing another thing where we need to improve on. We will definitely look into these things and thanks for all the valuable feedbacks and improvement ideas. Please share with us any feedback you have as well based on your evaluation.
🎅 1
thankyou 1
b
Hey @big-coat-53708, thanks for sharing!! Definitely agreed we should query Hive metastore, and even better, leverage multi-processing (for both extraction and publish). We at my company are internally working on a
Presto on Hive
plugin now to address the scalability issue and to support Presto views. We can share with the community once it’s done (cc @dazzling-judge-80093 we can collaborate on this if your team is on it too). To ingest Hive tables from Hive Metastore, i see there are 2 ways to achieve it among data catalogs.
1. Query HMS (Hive Metastore Service) via Thrift
2. Directly query Metastore DB
Alation adopts the former one, empirically it takes ~2h to ingest around 80K tables and 8K Presto views from Hive for Alation. Amundsen adopts the latter one (link). Each has pros and cons. But i would favor second one for best performance gain. We are learning the implementation from Amundsen for this. 😄
🎅 1
b
Hi @dazzling-judge-80093 @boundless-student-48844, thanks for the feedback! We also have a
trino + metastore
environment in our company. We extracted the views with the presto_view_extractor, it’s basically nothing different from the Hive extractor you pasted above. Just sharing in case you don’t know about it 😃 I don’t know much about the Hive plugin, but are you trying to implement
stateful ingestion
? Actually, the stateful ingestion or incremental pulling is the most needed feature for us. We will fully migrate to DataHub if it is supported for Hive Metastore. I believe all these latency won’t be a problem anymore since it will almost be realtime if we have stateful ingestion right? I know this feature has always been on the roadmap but does anyone know what is the latest status of it 🥲
b
@big-coat-53708, that’s really helpful! Thanks a lot for sharing it. We’ll borrow similar logic even for views in our
Presto on Hive
plugin on DataHub! As for stateful ingestion, it is released in 0.8.20. You can check out this month’s all hands - @/Shirshanka Das has a mention of it. But that’s more to address cases when an entity is removed from source. It won’t help to achieve real time ingestion or reduce time taken to ingest. To achieve real time, you would need to push directly from Hive to DataHub. A Hive hook would be required - but unfortunately there’s none from the community yet as I know of. Our company has some use cases that require real time metadata for Hive, we can collaborate in Q2/Q3 2022 if that fits your timeline too
b
@boundless-student-48844 Thanks for the clarification, I thought the stateful ingestion release is about the full push-based ingestion 😅. About the Hive hook and almost realtime ingestion, I think it’s also on the roadmap? @mammoth-bear-12532 mentioned about it here. Does anyone know what is the progress of this? Thanks 🙏
👍 1