Batch ingesting into multi-tenant collection

Description

We are having issues ingesting our big initial load of objects to the Weaviate instance.

Our collection is configured with a dynamic index (10_000 threshold) and auto creation/activation

We have 57 million objects distributed across 26k tenants. Only ~400 tenants have more than the 10_000 objects which results in hnsw index.

Our ingestion pipeline is a spark pipeline which, for each worker uses the .stream() api as recommended by the documentation. From what I can tell the ingest sends objects in no particular order, which then activates a bunch of tenants. This increases our memory usage, because the Weaviate keeps all active tenants in memory.

This then results in our Weaviate gets stuck in an OOM loop caused by the amount of active tenants.

Is there some approach we have misunderstood or could activation of flat indexes be circumvented during ingest on the server side?

Thank you in advance

Anton

Server Setup Information

  • Weaviate Server Version: 1.38.2
  • Deployment Method: K8
  • Multi Node? Number of Running Nodes: Yes. Three weaviate instances.
  • Client Language and Version: Python 4.22.0
  • Multitenancy?: Yes

Hi @bahtman !! Glad to see you back :slight_smile:

Is LIMIT_RESOURCES=true (or GOMEMLIMIT) set? — without it .stream()'s back-pressure is a no-op and that alone can explain the OOM loop.

Then repartition the Spark job by tenant and deactivate tenants as you finish them, so the hot set stays bounded.

Upgrading to 1.38.14 is worth doing regardless, especially if you’re running RF > 1.

1.38.14 is the current 1.38 patch, and several fixes in that range are directly on your path:

  • Replica resolution no longer implicitly activates tenants. This one matters a lot at RF > 1 — activating a tenant on one node could pull it hot on its replicas too, multiplying your resident tenant count by the replication factor.
  • Reduced allocations in tenant-activity tracking, and the tenant map is now cleared on an activity drop.
  • BatchStream goroutine leaks and shutdown races fixed — relevant when a stream dies mid-ingest and Spark retries it.
  • The memory monitor now refreshes while shards are read-only, so a node that hits the read-only guard can recover instead of staying stuck.

Let me know if this helps!

THanks!

Hi @DudaNogueira - Great to be back :wink:

Thank you so much for the quick reply.

We have LIMIT_RESOURCES=true but not GOMEMLIMIT set.

Wrt. repartition this is the approach we have had prior to the new .stream() api. We can easily revert to that, but a repartition (especially on initial load) is very expensive with that amount of tenants. Would be great if stream() has some built-in “auto-deactivation” and/or writing directly to disk on flat indexes - sorry if that makes no sense given the architecture :slight_smile:

For upgrading, the newest helmchart is 1.38.2. Is it sufficient to just bump the tag on our deploy?

Hi!

You can always push new weaviate verisons on top of those helm charts! :slight_smile: