每日写入2000万文档的ES集群,批量查询耗时30分钟是否正常?
Hey there, let's break down why your full-document scroll is dragging on for 30+ minutes even after tweaking size and using constant_score. With your setup—20M daily docs, 5 high-spec nodes, and 2-5 primary/replica shards—there are several targeted optimizations you can try to slash that runtime:
Your shard count might be holding back query performance:
- Align primary shards with node count: Set your index to 5 primary shards (matching your 5-node cluster). This lets each node handle exactly one primary shard, eliminating inter-node coordination overhead and maximizing parallel processing for scroll queries. Keep replicas at 2-3 (any more adds unnecessary IO/network load).
- Check shard size: Aim for shards between 5-15GB for query-heavy workloads (10-30GB is general guidance). If your shards are larger than 15GB (e.g., 20M docs at 1KB each = 20GB total, 2 shards would be 10GB each—this is okay, but 5 shards split it to 4GB each, which is more efficient for scrolling). If shards are too big, plan a reindex to split them (use a zero-downtime workflow if possible).
The biggest win for full-scroll speed in ES 6.x is slice scroll, which splits your query into parallel sub-queries that run across shards simultaneously:
- Split the scroll into N slices (match your primary shard count for optimal parallelism, e.g., 5 slices). Each slice processes a portion of the index independently.
- Example query for slice 0 of 5:
GET /your_index/_search?scroll=1m { "size": 5000, "slice": { "id": 0, "max": 5 }, "_source": ["critical_field_1", "critical_field_2"], // Only fetch needed fields! "query": { "constant_score": { "filter": { "match_all": {} } } } }
- Run 5 parallel requests (one for each slice ID 0-4), process each scroll batch independently, and merge results at the end. This can cut runtime by 70-90% depending on your cluster.
You’ve adjusted size, but there’s more to tweak:
- Avoid oversized
sizevalues: Stick to 1000-5000 per batch. Larger batches force more data transfer and processing per request, slowing things down. Test different sizes to find your sweet spot. - Trim
_sourceaggressively: Never fetch the full document unless you absolutely need it. Explicitly list only the fields you require—this reduces network bandwidth and memory usage drastically. - Manage scroll timeout wisely: Set
scroll=1minstead of a longer window (ES auto-extends the timeout as long as you keep scrolling). Always delete scroll IDs when done withDELETE /_search/scroll/{scroll_id}to free up cluster resources. - Force primary shard queries: Add
"preference": "_primary"to your query to avoid hitting replicas, which cuts inter-node network traffic. If primary shards are under heavy load, use"_local"to prioritize local shards instead.
Even high-spec hardware can hit bottlenecks if misconfigured:
- Check heap memory allocation: Allocate 50% of physical memory to ES heap, capped at 31GB (since JVMs struggle with heaps larger than 32GB). Leave the rest for Lucene’s file system cache—this is critical for fast query performance, as Lucene caches frequently accessed data here.
- Verify disk speed: If you’re using mechanical HDDs, switch to SSDs. Scroll queries are read-heavy, and SSDs can boost read throughput by 10-100x.
- Adjust search thread pool: The default search thread pool size is
int((available_processors * 3)/2) + 1. For high-core nodes, you can bump this slightly (but don’t exceed 2x CPU cores) to handle more parallel scroll slices. - Tweak refresh interval: If your index has a short
refresh_interval(e.g., 1s), Lucene generates too many small segments, slowing queries. Set it to 30s-1m during peak writes, or add"refresh": falseto your scroll query to avoid waiting for uncommitted segments.
Your daily 20M doc writes are competing for the same cluster resources as your scroll query. Run the full scroll during low-traffic windows to avoid IO/CPU contention—this alone can shave significant time off your runtime.
Start with slice scroll first—it’s the most impactful change for parallelizing your full index fetch. Then adjust your shard count and resource settings to support that parallelism. You should see a dramatic reduction in runtime once these changes are in place.
内容的提问来源于stack exchange,提问作者Duncan Krebs

