From 30437f24cd3470f60e7bc8501278d0babcb819e6 Mon Sep 17 00:00:00 2001 From: SungJin1212 Date: Tue, 6 Oct 2026 20:33:55 +0900 Subject: [PATCH 1/3] add support for honoring projection hints in Parquet mode bucket store Signed-off-by: SungJin1212 --- CHANGELOG.md | 1 + docs/blocks-storage/querier.md | 7 + docs/blocks-storage/store-gateway.md | 7 + docs/configuration/config-file-reference.md | 7 + docs/configuration/v1-guarantees.md | 2 + integration/parquet_store_gateway_test.go | 240 +++++++++ pkg/querier/querier.go | 12 +- pkg/querier/querier_test.go | 126 ++++- pkg/storage/tsdb/config.go | 4 +- pkg/storegateway/bucket_stores_bench_test.go | 2 +- pkg/storegateway/parquet_bucket_store.go | 48 +- .../parquet_bucket_store_bench_test.go | 184 ++++++- pkg/storegateway/parquet_bucket_store_test.go | 464 ++++++++++++++++++ pkg/storegateway/parquet_bucket_stores.go | 41 +- schemas/cortex-config-schema.json | 6 + 15 files changed, 1114 insertions(+), 37 deletions(-) create mode 100644 integration/parquet_store_gateway_test.go create mode 100644 pkg/storegateway/parquet_bucket_store_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 5087ac84c1b..a20f40b55ca 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ * [BUGFIX] Compactor: Fix the final cleanup of a tenant marked for deletion being reported as failed on object stores that return an error when deleting a missing object (GCS, Azure, Swift, OCI). #7861 * [FEATURE] Ruler: Add experimental support for federated rule groups. A rule group listing tenants in its `source_tenants` field is evaluated against those tenants while the resulting series and alerts are written to the tenant owning the rule group. Enabled with `-ruler.enable-federated-rules` (requires `-tenant-federation.enabled`), and restricted to selected tenants with `-ruler.allowed-federated-tenants` and `-ruler.disallowed-federated-tenants`. #7828 +* [FEATURE] StoreGateway: Add a flag `-blocks-storage.bucket-store.honor-projection-hints`. If enabled, Store Gateway in Parquet mode will honor projection hints and only materialize requested labels. #7206 ## 1.22.0 in progress * [CHANGE] Ruler: Remove the deprecated `-ruler.evaluation-delay-duration` flag and its `ruler_evaluation_delay_duration` per-tenant limit. Use `-ruler.query-offset` / `ruler_query_offset`, which no longer takes the higher of the two values. Cortex decodes the runtime config strictly, so a leftover `ruler_evaluation_delay_duration` override makes the runtime config fail to load: Cortex **exits at startup** (`module failed`, `module=runtime-config`), and on an already-running process every reload fails, pinning the last good overrides and dropping `cortex_runtime_config_last_reload_successful` to 0. Run `grep -r ruler_evaluation_delay_duration` over your runtime configs before upgrading. #7792 diff --git a/docs/blocks-storage/querier.md b/docs/blocks-storage/querier.md index 166b8f43635..5796d036b73 100644 --- a/docs/blocks-storage/querier.md +++ b/docs/blocks-storage/querier.md @@ -2166,6 +2166,13 @@ blocks_storage: # CLI flag: -blocks-storage.bucket-store.parquet-shard-cache-ttl [parquet_shard_cache_ttl: | default = 24h] + # [Experimental] If enabled, Store Gateway will honor projection hints and + # only materialize requested labels. It only takes effect when + # `-blocks-storage.bucket-store.bucket-store-type` is parquet and + # `-querier.honor-projection-hints` is enabled. + # CLI flag: -blocks-storage.bucket-store.honor-projection-hints + [honor_projection_hints: | default = false] + # Maximum number of concurrent goroutines per query applied at each level of # parquet processing: shard querying, row group processing, and column # materialization. Note: this limit is applied independently at each level, diff --git a/docs/blocks-storage/store-gateway.md b/docs/blocks-storage/store-gateway.md index 8c43bcb4653..4ff4ef2836f 100644 --- a/docs/blocks-storage/store-gateway.md +++ b/docs/blocks-storage/store-gateway.md @@ -2219,6 +2219,13 @@ blocks_storage: # CLI flag: -blocks-storage.bucket-store.parquet-shard-cache-ttl [parquet_shard_cache_ttl: | default = 24h] + # [Experimental] If enabled, Store Gateway will honor projection hints and + # only materialize requested labels. It only takes effect when + # `-blocks-storage.bucket-store.bucket-store-type` is parquet and + # `-querier.honor-projection-hints` is enabled. + # CLI flag: -blocks-storage.bucket-store.honor-projection-hints + [honor_projection_hints: | default = false] + # Maximum number of concurrent goroutines per query applied at each level of # parquet processing: shard querying, row group processing, and column # materialization. Note: this limit is applied independently at each level, diff --git a/docs/configuration/config-file-reference.md b/docs/configuration/config-file-reference.md index 6fa991fa3c4..9f69f5d6201 100644 --- a/docs/configuration/config-file-reference.md +++ b/docs/configuration/config-file-reference.md @@ -2863,6 +2863,13 @@ bucket_store: # CLI flag: -blocks-storage.bucket-store.parquet-shard-cache-ttl [parquet_shard_cache_ttl: | default = 24h] + # [Experimental] If enabled, Store Gateway will honor projection hints and + # only materialize requested labels. It only takes effect when + # `-blocks-storage.bucket-store.bucket-store-type` is parquet and + # `-querier.honor-projection-hints` is enabled. + # CLI flag: -blocks-storage.bucket-store.honor-projection-hints + [honor_projection_hints: | default = false] + # Maximum number of concurrent goroutines per query applied at each level of # parquet processing: shard querying, row group processing, and column # materialization. Note: this limit is applied independently at each level, so diff --git a/docs/configuration/v1-guarantees.md b/docs/configuration/v1-guarantees.md index b7a84922803..01bc029fe2f 100644 --- a/docs/configuration/v1-guarantees.md +++ b/docs/configuration/v1-guarantees.md @@ -97,6 +97,8 @@ Currently experimental features are: - Store Gateway Zone Stable Shuffle Sharding - `-store-gateway.sharding-ring.zone-stable-shuffle-sharding` CLI flag - `zone_stable_shuffle_sharding` (boolean) field in config file +- Store Gateway HonorProjectionHints in Parquet Mode + - `-blocks-storage.bucket-store.honor-projection-hints` CLI flag - Basic Lifecycler (Storegateway, Alertmanager, Ruler) Final Sleep on shutdown, which tells the pod wait before shutdown, allowing a delay to propagate ring changes. - `-ruler.ring.final-sleep` (duration) CLI flag - `store-gateway.sharding-ring.final-sleep` (duration) CLI flag diff --git a/integration/parquet_store_gateway_test.go b/integration/parquet_store_gateway_test.go new file mode 100644 index 00000000000..455d138a5cd --- /dev/null +++ b/integration/parquet_store_gateway_test.go @@ -0,0 +1,240 @@ +//go:build integration_querier + +package integration + +import ( + "context" + "fmt" + "math/rand" + "path/filepath" + "slices" + "testing" + "time" + + "github.com/prometheus/common/model" + "github.com/prometheus/prometheus/model/labels" + "github.com/stretchr/testify/require" + "github.com/thanos-io/thanos/pkg/block" + "github.com/thanos-io/thanos/pkg/block/metadata" + + "github.com/cortexproject/cortex/integration/e2e" + e2ecache "github.com/cortexproject/cortex/integration/e2e/cache" + e2edb "github.com/cortexproject/cortex/integration/e2e/db" + "github.com/cortexproject/cortex/integration/e2ecortex" + "github.com/cortexproject/cortex/pkg/storage/bucket" + "github.com/cortexproject/cortex/pkg/util/log" + cortex_testutil "github.com/cortexproject/cortex/pkg/util/test" +) + +func TestParquetBucketStore_ProjectionHint(t *testing.T) { + s, err := e2e.NewScenario(networkName) + require.NoError(t, err) + defer s.Close() + + consul := e2edb.NewConsulWithName("consul") + minio := e2edb.NewMinio(9000, bucketName) + memcached := e2ecache.NewMemcached() + require.NoError(t, s.StartAndWaitReady(consul, minio, memcached)) + + // Define configuration flags. + flags := BlocksStorageFlags() + flags = mergeFlags(flags, map[string]string{ + // Enable Thanos engine and projection optimization. + "-querier.thanos-engine": "true", + "-querier.optimizers": "projection", + + // enable honor-projection-hints querier and store gateway + "-querier.honor-projection-hints": "true", + "-blocks-storage.bucket-store.honor-projection-hints": "true", + // enable Store Gateway Parquet mode + "-blocks-storage.bucket-store.bucket-store-type": "parquet", + + // Set query-ingesters-within to 1h so queries older than 1h don't hit ingesters + "-limits.query-ingesters-within": "1h", + + // Configure Parquet Converter + "-parquet-converter.enabled": "true", + "-parquet-converter.conversion-interval": "1s", + "-parquet-converter.ring.consul.hostname": consul.NetworkHTTPEndpoint(), + "-compactor.block-ranges": "1ms,12h", + // Enable cache + "-blocks-storage.bucket-store.parquet-labels-cache.backend": "inmemory,memcached", + "-blocks-storage.bucket-store.parquet-labels-cache.memcached.addresses": "dns+" + memcached.NetworkEndpoint(e2ecache.MemcachedPort), + "-blocks-storage.bucket-store.sync-interval": "1s", + + // Compactor + "-compactor.cleanup-interval": "1s", // to update bucket index quickly + }) + + // Store Gateway + storeGateway := e2ecortex.NewStoreGateway("store-gateway", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), flags, "") + require.NoError(t, s.StartAndWaitReady(storeGateway)) + + // Parquet Converter + parquetConverter := e2ecortex.NewParquetConverter("parquet-converter", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), flags, "") + require.NoError(t, s.StartAndWaitReady(parquetConverter)) + + // Querier (honor projection hints enabled on the querier side) + querier := e2ecortex.NewQuerier("querier", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), mergeFlags(flags, map[string]string{ + "-querier.store-gateway-addresses": storeGateway.NetworkGRPCEndpoint(), + }), "") + require.NoError(t, s.StartAndWaitReady(querier)) + + // Querier with honor-projection-hints disabled on the querier side. The querier drops the + // projection hints, so the Store Gateway materializes all labels. + querierHonorOff := e2ecortex.NewQuerier("querier-honor-off", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), mergeFlags(flags, map[string]string{ + "-querier.store-gateway-addresses": storeGateway.NetworkGRPCEndpoint(), + "-querier.honor-projection-hints": "false", + }), "") + require.NoError(t, s.StartAndWaitReady(querierHonorOff)) + + require.NoError(t, querier.WaitSumMetrics(e2e.Equals(512), "cortex_ring_tokens_total")) + require.NoError(t, querierHonorOff.WaitSumMetrics(e2e.Equals(512), "cortex_ring_tokens_total")) + require.NoError(t, storeGateway.WaitSumMetrics(e2e.Equals(512), "cortex_ring_tokens_total")) + + // Create block + now := time.Now() + // Time range: [Now - 24h] to [Now - 20h] + start := now.Add(-24 * time.Hour) + end := now.Add(-20 * time.Hour) + + ctx := context.Background() + + rnd := rand.New(rand.NewSource(time.Now().Unix())) + dir := filepath.Join(s.SharedDir(), "data") + scrapeInterval := time.Minute + statusCodes := []string{"200", "400", "404", "500", "502"} + methods := []string{"GET", "POST", "PUT", "DELETE"} + + numSeries := 10 + numSamples := 100 + + lbls := make([]labels.Labels, 0, numSeries) + for i := 0; i < numSeries; i++ { + lbls = append(lbls, labels.FromStrings( + labels.MetricName, "http_requests_total", + "job", "api-server", + "instance", fmt.Sprintf("instance-%d", i), + "status_code", statusCodes[i%len(statusCodes)], + "method", methods[i%len(methods)], + "path", fmt.Sprintf("/api/v1/endpoint%d", i%3), + "cluster", "test-cluster", + )) + } + + id, err := e2e.CreateBlock(ctx, rnd, dir, lbls, numSamples, start.UnixMilli(), end.UnixMilli(), scrapeInterval.Milliseconds(), 10) + require.NoError(t, err) + + storage, err := e2ecortex.NewS3ClientForMinio(minio, bucketName) + require.NoError(t, err) + bkt := bucket.NewUserBucketClient("user-1", storage.GetBucket(), nil) + + // Upload TSDB Block + require.NoError(t, block.Upload(ctx, log.Logger, bkt, filepath.Join(dir, id.String()), metadata.NoneFunc)) + + // Start compactor to create the bucket index. + compactor := e2ecortex.NewCompactor("compactor", consul.NetworkHTTPEndpoint(), flags, "") + require.NoError(t, s.StartAndWaitReady(compactor)) + + // Wait until parquet converter convert block + require.NoError(t, parquetConverter.WaitSumMetricsWithOptions(e2e.Equals(float64(1)), []string{"cortex_parquet_converter_blocks_converted_total"}, e2e.WaitMissingMetrics)) + + // Create clients for both queriers (honor enabled / disabled on the querier side). + c, err := e2ecortex.NewClient("", querier.HTTPEndpoint(), "", "", "user-1") + require.NoError(t, err) + cHonorOff, err := e2ecortex.NewClient("", querierHonorOff.HTTPEndpoint(), "", "", "user-1") + require.NoError(t, err) + + cortex_testutil.Poll(t, 60*time.Second, true, func() interface{} { + labelSets, err := c.Series([]string{`{job="api-server"}`}, start, end) + if err != nil { + t.Logf("Series query failed: %v", err) + return false + } + return len(labelSets) > 0 + }) + + testCases := []struct { + name string + query string + expectedLabels []string // query result should contain these labels + }{ + { + name: "vector selector query", + query: `http_requests_total`, + expectedLabels: []string{ + "__name__", "job", "instance", "status_code", "method", "path", "cluster", + }, + }, + { + name: "simple_sum_by_job", + query: `sum by (job) (http_requests_total)`, + expectedLabels: []string{"job"}, + }, + { + name: "rate_with_aggregation", + query: `sum by (method) (rate(http_requests_total[5m]))`, + expectedLabels: []string{"method"}, + }, + { + name: "multiple_grouping_labels", + query: `sum by (job, status_code) (http_requests_total)`, + expectedLabels: []string{"job", "status_code"}, + }, + { + name: "aggregation without query", + query: `sum without (instance, method) (http_requests_total)`, + expectedLabels: []string{"job", "status_code", "path", "cluster"}, + }, + } + // Run the same assertions against both queriers. The querier with honor-projection-hints + // disabled must produce identical results even though the Store Gateway does not project. + clients := []struct { + name string + client *e2ecortex.Client + }{ + {name: "querier_honor_on", client: c}, + {name: "querier_honor_off", client: cHonorOff}, + } + + for _, qc := range clients { + t.Run(qc.name, func(t *testing.T) { + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + t.Logf("Testing: %s", tc.query) + + // Execute instant query + result, err := qc.client.Query(tc.query, end) + require.NoError(t, err) + require.NotNil(t, result) + + // Verify we got results + vector, ok := result.(model.Vector) + require.True(t, ok, "result should be a vector") + require.NotEmpty(t, vector, "query should return results") + + for _, sample := range vector { + actualLabels := make(map[string]struct{}) + for label := range sample.Metric { + actualLabels[string(label)] = struct{}{} + } + + // Check that all expected labels are present + for _, expectedLabel := range tc.expectedLabels { + _, ok := actualLabels[expectedLabel] + require.True(t, ok, + "series should have %s label", expectedLabel) + } + + // Check that no unexpected labels are present + for lbl := range actualLabels { + if !slices.Contains(tc.expectedLabels, lbl) { + require.Fail(t, "series should not have unexpected label: %s", lbl) + } + } + } + }) + } + }) + } +} diff --git a/pkg/querier/querier.go b/pkg/querier/querier.go index dc9edd6dc56..44cd05e3817 100644 --- a/pkg/querier/querier.go +++ b/pkg/querier/querier.go @@ -529,13 +529,11 @@ func (q querier) Select(ctx context.Context, sortSeries bool, sp *storage.Select } } - // Reset projection hints if querying ingesters or projection is not included. - // Projection can only be applied when not querying mixed sources (ingester + store). - if q.honorProjectionHints { - if !sp.ProjectionInclude || q.distributor.UseQueryable(q.now, userID, mint, maxt) { - sp.ProjectionLabels = nil - sp.ProjectionInclude = false - } + // Reset projection hints unless the querier honors them, projection is included and ingesters are not queried. + // Ingesters always return full label sets, so projected store series could not be merged with them. + if !q.honorProjectionHints || !sp.ProjectionInclude || q.distributor.UseQueryable(q.now, userID, mint, maxt) { + sp.ProjectionLabels = nil + sp.ProjectionInclude = false } if len(queriers) == 1 { diff --git a/pkg/querier/querier_test.go b/pkg/querier/querier_test.go index 53f11c8fa08..66f81a0883b 100644 --- a/pkg/querier/querier_test.go +++ b/pkg/querier/querier_test.go @@ -11,7 +11,9 @@ import ( "time" "github.com/go-kit/log" + "github.com/oklog/ulid/v2" "github.com/pkg/errors" + "github.com/prometheus-community/parquet-common/schema" "github.com/prometheus/client_golang/prometheus" promutil "github.com/prometheus/client_golang/prometheus/testutil" "github.com/prometheus/common/model" @@ -29,7 +31,9 @@ import ( "github.com/stretchr/testify/require" "github.com/thanos-io/promql-engine/engine" "github.com/thanos-io/promql-engine/logicalplan" + "github.com/thanos-io/thanos/pkg/store/storepb" "github.com/weaveworks/common/user" + "google.golang.org/grpc" "github.com/cortexproject/cortex/pkg/chunk" promchunk "github.com/cortexproject/cortex/pkg/chunk/encoding" @@ -38,6 +42,8 @@ import ( cortexparser "github.com/cortexproject/cortex/pkg/parser" "github.com/cortexproject/cortex/pkg/querier/batch" "github.com/cortexproject/cortex/pkg/querier/series" + "github.com/cortexproject/cortex/pkg/storage/tsdb/bucketindex" + "github.com/cortexproject/cortex/pkg/storegateway/storegatewaypb" "github.com/cortexproject/cortex/pkg/util" "github.com/cortexproject/cortex/pkg/util/chunkcompat" "github.com/cortexproject/cortex/pkg/util/flagext" @@ -1839,21 +1845,29 @@ func TestQuerier_ProjectionHints(t *testing.T) { expectedProjectionInclude: false, expectedProjectionLabels: nil, }, - "projection not modified: honor disabled, projection included, querying ingesters": { + "projection reset: honor disabled, projection included, querying ingesters": { honorProjectionHints: false, inputProjectionInclude: true, inputProjectionLabels: []string{"__name__", "job"}, queryIngesters: true, - expectedProjectionInclude: true, - expectedProjectionLabels: []string{"__name__", "job"}, + expectedProjectionInclude: false, + expectedProjectionLabels: nil, + }, + "projection reset: honor disabled, projection included, not querying ingesters": { + honorProjectionHints: false, + inputProjectionInclude: true, + inputProjectionLabels: []string{"__name__", "job"}, + queryIngesters: false, + expectedProjectionInclude: false, + expectedProjectionLabels: nil, }, - "projection not modified: honor disabled, projection not included": { + "projection reset: honor disabled, projection not included": { honorProjectionHints: false, inputProjectionInclude: false, inputProjectionLabels: []string{"__name__", "job"}, queryIngesters: false, expectedProjectionInclude: false, - expectedProjectionLabels: []string{"__name__", "job"}, + expectedProjectionLabels: nil, }, } @@ -1932,6 +1946,108 @@ func TestQuerier_ProjectionHints(t *testing.T) { } } +// projectingStoreGatewayClientMock mimics a parquet store gateway running with +// -blocks-storage.bucket-store.honor-projection-hints=true. +type projectingStoreGatewayClientMock struct { + *storeGatewayClientMock + series labels.Labels + samples []cortexpb.Sample + blockID ulid.ULID +} + +func (m *projectingStoreGatewayClientMock) Series(ctx context.Context, in *storepb.SeriesRequest, opts ...grpc.CallOption) (storegatewaypb.StoreGateway_SeriesClient, error) { + lbls := m.series + if in.QueryHints != nil && in.QueryHints.ProjectionInclude { + b := labels.NewBuilder(labels.EmptyLabels()) + for _, name := range in.QueryHints.ProjectionLabels { + b.Set(name, m.series.Get(name)) + } + b.Set(schema.SeriesHashColumn, strconv.FormatUint(m.series.Hash(), 10)) + lbls = b.Labels() + } + m.mockedSeriesResponses = []*storepb.SeriesResponse{ + mockSeriesResponse(lbls, m.samples, nil, nil), + mockHintsResponse(m.blockID), + } + return m.storeGatewayClientMock.Series(ctx, in, opts...) +} + +// Ensures a projecting store gateway never returns projected series for a query that +// also hits ingesters, regardless of the querier's honor-projection-hints setting. +func TestQuerier_ProjectionHints_IngestersAndProjectingStoreGateway(t *testing.T) { + t.Parallel() + const userID = "user-1" + end := time.Now() + start := end.Add(-time.Hour) + minT, maxT := util.TimeToMillis(start), util.TimeToMillis(end) + block1 := ulid.MustNew(1, nil) + + seriesLabels := labels.FromStrings(labels.MetricName, "test", "job", "a", "instance", "i1") + samples := []cortexpb.Sample{{TimestampMs: minT, Value: 1}, {TimestampMs: minT + 1000, Value: 2}} + + for _, honorProjectionHints := range []bool{true, false} { + t.Run(fmt.Sprintf("querier honor-projection-hints=%t", honorProjectionHints), func(t *testing.T) { + t.Parallel() + var cfg Config + flagext.DefaultValues(&cfg) + cfg.ActiveQueryTrackerDir = "" + cfg.HonorProjectionHints = honorProjectionHints + + // Ingesters hold the same samples, returned with the full label set. + matrix := model.Matrix{{Metric: util.LabelsToMetric(seriesLabels)}} + for _, s := range samples { + matrix[0].Values = append(matrix[0].Values, model.SamplePair{Timestamp: model.Time(s.TimestampMs), Value: model.SampleValue(s.Value)}) + } + ingesterQueryable := UseAlwaysQueryable(storage.QueryableFunc(func(_, _ int64) (storage.Querier, error) { + return mockQuerier{matrix: matrix}, nil + })) + + sgClient := &projectingStoreGatewayClientMock{ + storeGatewayClientMock: &storeGatewayClientMock{remoteAddr: "1.1.1.1"}, + series: seriesLabels, + samples: samples, + blockID: block1, + } + storeQueryable := UseAlwaysQueryable(storage.QueryableFunc(func(mint, maxt int64) (storage.Querier, error) { + finder := &blocksFinderMock{} + finder.On("GetBlocks", mock.Anything, userID, mock.Anything, mock.Anything, mock.Anything).Return(bucketindex.Blocks{&bucketindex.Block{ID: block1}}, map[ulid.ULID]*bucketindex.BlockDeletionMark(nil), nil) + return &blocksStoreQuerier{ + minT: mint, + maxT: maxt, + finder: finder, + stores: &blocksStoreSetMock{mockedResponses: []any{map[BlocksStoreClient][]ulid.ULID{sgClient: {block1}}}}, + consistency: NewBlocksConsistencyChecker(0, 0, log.NewNopLogger(), nil), + logger: log.NewNopLogger(), + metrics: newBlocksStoreQueryableMetrics(prometheus.NewPedanticRegistry()), + limits: &blocksStoreLimitsMock{}, + + storeGatewayConsistencyCheckMaxAttempts: 1, + }, nil + })) + + overrides := validation.NewOverrides(DefaultLimitsConfig(), nil) + queryable := NewQueryable(ingesterQueryable, []QueryableWithFilter{storeQueryable}, cfg, overrides, nil, log.NewNopLogger(), nil) + q, err := queryable.Querier(minT, maxT) + require.NoError(t, err) + + // Hints as produced by the Thanos engine projection optimizer for `sum by (job) (test)`. + hints := &storage.SelectHints{Start: minT, End: maxT, ProjectionInclude: true, ProjectionLabels: []string{"job"}} + ctx := limiter.AddQueryLimiterToContext(user.InjectOrgID(context.Background(), userID), limiter.NewQueryLimiter(0, 0, 0, 0)) + set := q.Select(ctx, true, hints, labels.MustNewMatcher(labels.MatchEqual, labels.MetricName, "test")) + + var got []labels.Labels + for set.Next() { + got = append(got, set.At().Labels()) + } + require.NoError(t, set.Err()) + + // The same series is held by both ingesters and the store gateway, so it must be + // merged into one; otherwise overlapping samples are counted twice by aggregations. + assert.Equal(t, []labels.Labels{seriesLabels}, got) + }) + } +} + func TestQuerier_ResourceBasedLimiter(t *testing.T) { cfg := Config{} flagext.DefaultValues(&cfg) diff --git a/pkg/storage/tsdb/config.go b/pkg/storage/tsdb/config.go index db51ebc5765..7bbcddcc66c 100644 --- a/pkg/storage/tsdb/config.go +++ b/pkg/storage/tsdb/config.go @@ -334,7 +334,8 @@ type BucketStoreConfig struct { // Token bucket configs TokenBucketBytesLimiter TokenBucketBytesLimiterConfig `yaml:"token_bucket_bytes_limiter"` // Parquet shard cache config - ParquetShardCache parquetutil.CacheConfig `yaml:",inline"` + ParquetShardCache parquetutil.CacheConfig `yaml:",inline"` + HonorProjectionHints bool `yaml:"honor_projection_hints"` // ParquetQueryConcurrency controls the maximum number of concurrent goroutines // per query at each level of parquet processing: shard querying, row group @@ -405,6 +406,7 @@ func (cfg *BucketStoreConfig) RegisterFlags(f *flag.FlagSet) { f.IntVar(&cfg.MatchersCacheMaxItems, "blocks-storage.bucket-store.matchers-cache-max-items", 0, "Maximum number of entries in the regex matchers cache. 0 to disable.") f.IntVar(&cfg.ParquetQueryConcurrency, "blocks-storage.bucket-store.parquet-query-concurrency", 4, "Maximum number of concurrent goroutines per query applied at each level of parquet processing: shard querying, row group processing, and column materialization. Note: this limit is applied independently at each level, so the total goroutines per query can grow multiplicatively (up to N^3 in the worst case).") cfg.ParquetShardCache.RegisterFlagsWithPrefix("blocks-storage.bucket-store.", f) + f.BoolVar(&cfg.HonorProjectionHints, "blocks-storage.bucket-store.honor-projection-hints", false, "[Experimental] If enabled, Store Gateway will honor projection hints and only materialize requested labels. It only takes effect when `-blocks-storage.bucket-store.bucket-store-type` is parquet and `-querier.honor-projection-hints` is enabled.") } // Validate the config. diff --git a/pkg/storegateway/bucket_stores_bench_test.go b/pkg/storegateway/bucket_stores_bench_test.go index 2dd9c2207f7..0fa3f003fe0 100644 --- a/pkg/storegateway/bucket_stores_bench_test.go +++ b/pkg/storegateway/bucket_stores_bench_test.go @@ -128,7 +128,7 @@ func benchmarkBatching(b *testing.B, client storepb.StoreClient, userID string, } } -func generateBenchmarkBlock(b *testing.B, storageDir, userID string, numSeries, numSamples int) { +func generateBenchmarkBlock(b testing.TB, storageDir, userID string, numSeries, numSamples int) { userDir := filepath.Join(storageDir, userID) if err := os.MkdirAll(userDir, os.ModePerm); err != nil { b.Fatal(err) diff --git a/pkg/storegateway/parquet_bucket_store.go b/pkg/storegateway/parquet_bucket_store.go index 246e68e9abc..19dda57b11d 100644 --- a/pkg/storegateway/parquet_bucket_store.go +++ b/pkg/storegateway/parquet_bucket_store.go @@ -3,6 +3,7 @@ package storegateway import ( "context" "fmt" + "slices" "strings" "sync" @@ -46,9 +47,10 @@ type parquetBucketStore struct { chunksDecoder *schema.PrometheusParquetChunksDecoder - matcherCache storecache.MatchersCache - parquetShardCache parquetutil.CacheInterface[parquet_storage.ParquetShard] - rowRangesCache search.RowRangesForConstraintsCache + matcherCache storecache.MatchersCache + parquetShardCache parquetutil.CacheInterface[parquet_storage.ParquetShard] + rowRangesCache search.RowRangesForConstraintsCache + honorProjectionHints bool shardCountsMu sync.Mutex // cachedShardCounts maps a block ID to its parquet shard count. The map is rebuilt @@ -232,6 +234,7 @@ func (p *parquetBucketStore) Series(req *storepb.SeriesRequest, seriesSrv storep return fmt.Errorf("failed to find parquet shards: %w", err) } + storageHints := p.buildSelectHints(req.QueryHints, shards, req.MinTime, req.MaxTime) seriesSet := make([]prom_storage.ChunkSeriesSet, len(shards)) errGroup, ctx := errgroup.WithContext(srv.Context()) errGroup.SetLimit(p.concurrency) @@ -245,7 +248,7 @@ func (p *parquetBucketStore) Series(req *storepb.SeriesRequest, seriesSrv storep }) } errGroup.Go(func() error { - ss, err := shard.Query(ctx, req.MinTime, req.MaxTime, req.SkipChunks, matchers) + ss, err := shard.Query(ctx, storageHints, req.SkipChunks, matchers) seriesSet[i] = ss return err }) @@ -416,3 +419,40 @@ func (p *parquetBucketStore) LabelValues(ctx context.Context, req *storepb.Label Hints: anyHints, }, nil } + +func (p *parquetBucketStore) buildSelectHints(queryHints *storepb.QueryHints, shards []*parquetBlock, minT, maxT int64) *prom_storage.SelectHints { + storageHints := &prom_storage.SelectHints{ + Start: minT, + End: maxT, + } + + if p.honorProjectionHints && queryHints != nil && queryHints.ProjectionInclude { + storageHints.ProjectionInclude = true + // Copy so appending the hash column never writes into the request's backing array. + lbls := make([]string, 0, len(queryHints.ProjectionLabels)+1) + storageHints.ProjectionLabels = append(lbls, queryHints.ProjectionLabels...) + + if !allParquetBlocksHaveHashColumn(shards) { + // Reset projection hints if not all parquet shards have the hash column (version < 2). + storageHints.ProjectionInclude = false + storageHints.ProjectionLabels = nil + } + + if storageHints.ProjectionInclude && !slices.Contains(storageHints.ProjectionLabels, schema.SeriesHashColumn) { + // Series hash column is always required for projection. + storageHints.ProjectionLabels = append(storageHints.ProjectionLabels, schema.SeriesHashColumn) + } + } + + return storageHints +} + +func allParquetBlocksHaveHashColumn(blocks []*parquetBlock) bool { + // TODO(Sungjin1212): Change it to read marker version + for _, b := range blocks { + if !b.hasHashColumn() { + return false + } + } + return true +} diff --git a/pkg/storegateway/parquet_bucket_store_bench_test.go b/pkg/storegateway/parquet_bucket_store_bench_test.go index 659c21b4c51..a8162fd8544 100644 --- a/pkg/storegateway/parquet_bucket_store_bench_test.go +++ b/pkg/storegateway/parquet_bucket_store_bench_test.go @@ -40,6 +40,176 @@ import ( "github.com/cortexproject/cortex/pkg/util/validation" ) +// BenchmarkParquetBucketStore_ProjectionHints compares query performance +// with and without projectionHints enabled. +func BenchmarkParquetBucketStore_ProjectionHints(b *testing.B) { + seriesNum := []int{100, 1000, 10000} + samplePerSeries := 100 + + // projectionLabels is the subset of labels requested via projection hints. + // The test data contains: __name__, idx, job, instance + projectionLabels := []string{"job"} + + for _, numSeries := range seriesNum { + b.Run(fmt.Sprintf("series_%d", numSeries), func(b *testing.B) { + ctx := context.Background() + tmpDir := b.TempDir() + storageDir := filepath.Join(tmpDir, "storage") + dataDir := filepath.Join(tmpDir, "data") + userID := "user-1" + + storageCfg := cortex_tsdb.BlocksStorageConfig{ + UsersScanner: users.UsersScannerConfig{ + Strategy: users.UserScanStrategyList, + UpdateInterval: time.Second, + }, + Bucket: bucket.Config{ + Backend: "filesystem", + Filesystem: filesystem.Config{ + Directory: storageDir, + }, + }, + BucketStore: cortex_tsdb.BucketStoreConfig{ + SyncDir: filepath.Join(tmpDir, "sync"), + BucketStoreType: "parquet", + BlockDiscoveryStrategy: string(cortex_tsdb.RecursiveDiscovery), + // Must be > 0, otherwise the per-query errgroup limit becomes 0 + // and errGroup.Go blocks forever, hanging the Series handler. + ParquetQueryConcurrency: 4, + }, + } + + bucketClient, err := bucket.NewClient(ctx, storageCfg.Bucket, nil, "test", log.NewNopLogger(), prometheus.NewRegistry()) + require.NoError(b, err) + + blockID := prepareParquetBlock(b, ctx, storageCfg, bucketClient, dataDir, userID, numSeries, samplePerSeries) + + startGRPCServer := func(honorProjectionHints bool) (storepb.StoreClient, func()) { + cfg := storageCfg + cfg.BucketStore.HonorProjectionHints = honorProjectionHints + + reg := prometheus.NewPedanticRegistry() + stores, err := NewBucketStores(cfg, NewNoShardingStrategy(log.NewNopLogger(), nil), objstore.WithNoopInstr(bucketClient), defaultLimitsOverrides(nil), mockLoggingLevel(), log.NewNopLogger(), reg) + require.NoError(b, err) + + listener, err := net.Listen("tcp", "localhost:0") + require.NoError(b, err) + + gRPCServer := grpc.NewServer( + grpc.StreamInterceptor(middleware.StreamServerUserHeaderInterceptor), + ) + storepb.RegisterStoreServer(gRPCServer, stores) + + go func() { + if err := gRPCServer.Serve(listener); err != nil && err != grpc.ErrServerStopped { + b.Error(err) + } + }() + + conn, err := grpc.NewClient(listener.Addr().String(), + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithDefaultCallOptions(grpc.MaxCallRecvMsgSize(math.MaxInt32)), + ) + require.NoError(b, err) + + return storepb.NewStoreClient(conn), func() { + _ = conn.Close() + gRPCServer.Stop() + } + } + + // Benchmark without projectionHints (full scan) + b.Run("without_projection_hints", func(b *testing.B) { + client, stop := startGRPCServer(false) + defer stop() + + b.ReportAllocs() + for b.Loop() { + benchmarkProjectionHints(b, client, userID, blockID, numSeries, nil) + } + }) + + // Benchmark with projectionHints enabled and projection labels specified + b.Run("with_projection_hints", func(b *testing.B) { + client, stop := startGRPCServer(true) + defer stop() + + b.ReportAllocs() + for b.Loop() { + benchmarkProjectionHints(b, client, userID, blockID, numSeries, projectionLabels) + } + }) + }) + } +} + +func benchmarkProjectionHints(b *testing.B, client storepb.StoreClient, userID, blockID string, expectedSeries int, projectionLabels []string) { + b.Helper() + + ctx := grpcMetadata.NewOutgoingContext(context.Background(), grpcMetadata.Pairs(cortex_tsdb.TenantIDExternalLabel, userID)) + ctx, err := user.InjectIntoGRPCRequest(user.InjectOrgID(ctx, userID)) + require.NoError(b, err) + + hintMatchers := []storepb.LabelMatcher{ + { + Type: storepb.LabelMatcher_RE, + Name: block.BlockIDLabel, + Value: blockID, + }, + } + + dataMatchers := []storepb.LabelMatcher{ + { + Type: storepb.LabelMatcher_RE, + Name: "__name__", + Value: ".+", + }, + } + + hints := &hintspb.SeriesRequestHints{ + BlockMatchers: hintMatchers, + } + hintsAny, err := types.MarshalAny(hints) + require.NoError(b, err) + + req := &storepb.SeriesRequest{ + MinTime: 0, + MaxTime: math.MaxInt64, + Matchers: dataMatchers, + ResponseBatchSize: 1000, + Hints: hintsAny, + } + + if len(projectionLabels) > 0 { + req.QueryHints = &storepb.QueryHints{ + ProjectionInclude: true, + ProjectionLabels: projectionLabels, + } + } + + stream, err := client.Series(ctx, req) + require.NoError(b, err) + + got := 0 + for { + resp, err := stream.Recv() + if err == io.EOF { + break + } + require.NoError(b, err) + + if series := resp.GetSeries(); series != nil { + got++ + } else if batch := resp.GetBatch(); batch != nil { + got += len(batch.Series) + } + } + + if got != expectedSeries { + b.Fatalf("expected %d series, got %d", expectedSeries, got) + } +} + func BenchmarkParquetBucketStore_SeriesBatch(b *testing.B) { seriesNum := []int{100, 1000, 10000, 100000} samplePerSeries := 100 @@ -69,6 +239,9 @@ func BenchmarkParquetBucketStore_SeriesBatch(b *testing.B) { SyncDir: filepath.Join(tmpDir, "sync"), BucketStoreType: "parquet", BlockDiscoveryStrategy: string(cortex_tsdb.RecursiveDiscovery), + // Must be > 0, otherwise the per-query errgroup limit becomes 0 + // and errGroup.Go blocks forever, hanging the Series handler. + ParquetQueryConcurrency: 4, }, } bucketClient, err := bucket.NewClient(context.Background(), storageCfg.Bucket, nil, "test", log.NewNopLogger(), prometheus.NewRegistry()) @@ -159,9 +332,10 @@ func BenchmarkParquetBucketStore_MultiShard(b *testing.B) { }, }, BucketStore: cortex_tsdb.BucketStoreConfig{ - SyncDir: filepath.Join(tmpDir, "sync"), - BucketStoreType: "parquet", - BlockDiscoveryStrategy: string(cortex_tsdb.RecursiveDiscovery), + SyncDir: filepath.Join(tmpDir, "sync"), + BucketStoreType: "parquet", + BlockDiscoveryStrategy: string(cortex_tsdb.RecursiveDiscovery), + ParquetQueryConcurrency: 4, }, } bucketClient, err := bucket.NewClient(ctx, storageCfg.Bucket, nil, "test", log.NewNopLogger(), prometheus.NewRegistry()) @@ -204,11 +378,11 @@ func BenchmarkParquetBucketStore_MultiShard(b *testing.B) { } } -func prepareParquetBlock(b *testing.B, ctx context.Context, storageCfg cortex_tsdb.BlocksStorageConfig, bkt objstore.InstrumentedBucket, dataDir, userID string, numSeries, numSamples int) string { +func prepareParquetBlock(b testing.TB, ctx context.Context, storageCfg cortex_tsdb.BlocksStorageConfig, bkt objstore.InstrumentedBucket, dataDir, userID string, numSeries, numSamples int) string { return prepareParquetBlockWithShards(b, ctx, storageCfg, bkt, dataDir, userID, numSeries, numSamples, math.MaxInt32, 1_000_000) } -func prepareParquetBlockWithShards(b *testing.B, ctx context.Context, storageCfg cortex_tsdb.BlocksStorageConfig, bkt objstore.InstrumentedBucket, dataDir, userID string, numSeries, numSamples, numRowGroups, maxRowsPerRowGroup int) string { +func prepareParquetBlockWithShards(b testing.TB, ctx context.Context, storageCfg cortex_tsdb.BlocksStorageConfig, bkt objstore.InstrumentedBucket, dataDir, userID string, numSeries, numSamples, numRowGroups, maxRowsPerRowGroup int) string { logger := log.NewNopLogger() reg := prometheus.NewRegistry() diff --git a/pkg/storegateway/parquet_bucket_store_test.go b/pkg/storegateway/parquet_bucket_store_test.go new file mode 100644 index 00000000000..2793b51df99 --- /dev/null +++ b/pkg/storegateway/parquet_bucket_store_test.go @@ -0,0 +1,464 @@ +package storegateway + +import ( + "bytes" + "context" + "io" + "math" + "net" + "path/filepath" + "testing" + "time" + + "github.com/go-kit/log" + "github.com/gogo/protobuf/types" + "github.com/parquet-go/parquet-go" + "github.com/prometheus-community/parquet-common/schema" + parquet_storage "github.com/prometheus-community/parquet-common/storage" + "github.com/prometheus/client_golang/prometheus" + prom_storage "github.com/prometheus/prometheus/storage" + "github.com/stretchr/testify/require" + "github.com/thanos-io/objstore" + "github.com/thanos-io/thanos/pkg/block" + "github.com/thanos-io/thanos/pkg/store/hintspb" + "github.com/thanos-io/thanos/pkg/store/labelpb" + "github.com/thanos-io/thanos/pkg/store/storepb" + "github.com/weaveworks/common/middleware" + "github.com/weaveworks/common/user" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + grpcMetadata "google.golang.org/grpc/metadata" + + "github.com/cortexproject/cortex/pkg/storage/bucket" + "github.com/cortexproject/cortex/pkg/storage/bucket/filesystem" + cortex_tsdb "github.com/cortexproject/cortex/pkg/storage/tsdb" + "github.com/cortexproject/cortex/pkg/util/users" +) + +type mockParquetFileView struct { + parquet_storage.ParquetFileView + file *parquet.File +} + +func (m *mockParquetFileView) Schema() *parquet.Schema { + return m.file.Schema() +} + +type mockShard struct { + parquet_storage.ParquetShard + fileView parquet_storage.ParquetFileView +} + +func (m *mockShard) LabelsFile() parquet_storage.ParquetFileView { + return m.fileView +} + +func createTestParquetBlock(t *testing.T, hasHash bool) *parquetBlock { + t.Helper() + + var buf bytes.Buffer + var err error + + if hasHash { + // v2 and higher blocks + type RowWithHash struct { + SeriesHash string `parquet:"s_series_hash"` + Label string `parquet:"l_job"` + } + + w := parquet.NewGenericWriter[RowWithHash](&buf) + _, err = w.Write([]RowWithHash{{SeriesHash: "hash1", Label: "node-1"}}) + require.NoError(t, err) + require.NoError(t, w.Close()) + } else { + // v1 block + type RowWithoutHash struct { + Label string `parquet:"l_job"` + } + + w := parquet.NewGenericWriter[RowWithoutHash](&buf) + _, err = w.Write([]RowWithoutHash{{Label: "node-1"}}) + require.NoError(t, err) + require.NoError(t, w.Close()) + } + + readBuf := bytes.NewReader(buf.Bytes()) + f, err := parquet.OpenFile(readBuf, readBuf.Size()) + require.NoError(t, err) + + return &parquetBlock{ + shard: &mockShard{ + fileView: &mockParquetFileView{file: f}, + }, + } +} + +func Test_AllParquetBlocksHaveHashColumn(t *testing.T) { + tests := []struct { + description string + setup func() []*parquetBlock + expected bool + }{ + { + description: "returns true when all blocks have hash column", + setup: func() []*parquetBlock { + return []*parquetBlock{ + createTestParquetBlock(t, true), + createTestParquetBlock(t, true), + createTestParquetBlock(t, true), + } + }, + expected: true, + }, + { + description: "returns false when mixed block versions exist", + setup: func() []*parquetBlock { + return []*parquetBlock{ + createTestParquetBlock(t, true), + createTestParquetBlock(t, false), + createTestParquetBlock(t, true), + } + }, + expected: false, + }, + { + description: "returns false when no blocks have hash column", + setup: func() []*parquetBlock { + return []*parquetBlock{ + createTestParquetBlock(t, false), + createTestParquetBlock(t, false), + } + }, + expected: false, + }, + { + description: "returns true for empty block list", + setup: func() []*parquetBlock { + return []*parquetBlock{} + }, + expected: true, + }, + } + + for _, tc := range tests { + t.Run(tc.description, func(t *testing.T) { + blocks := tc.setup() + actual := allParquetBlocksHaveHashColumn(blocks) + require.Equal(t, tc.expected, actual) + }) + } +} + +func TestParquetBucketStore_buildSelectHints(t *testing.T) { + const ( + minT = 1000 + maxT = 2000 + ) + + tests := []struct { + description string + honorProjectionHints bool + queryHints *storepb.QueryHints + shards []*parquetBlock + expectedHints *prom_storage.SelectHints + }{ + { + description: "honorProjectionHints=false, should ignore query hints", + honorProjectionHints: false, + queryHints: &storepb.QueryHints{ + ProjectionInclude: true, + ProjectionLabels: []string{"job"}, + }, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: false, + ProjectionLabels: nil, + }, + }, + { + description: "honorProjectionHints=true, V2 blocks, should enable projection and append hash column", + honorProjectionHints: true, + queryHints: &storepb.QueryHints{ + ProjectionInclude: true, + ProjectionLabels: []string{"job"}, + }, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + createTestParquetBlock(t, true), + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: true, + ProjectionLabels: []string{"job", schema.SeriesHashColumn}, + }, + }, + { + description: "honorProjectionHints=true, V2 blocks, should not duplicate hash column if already present", + honorProjectionHints: true, + queryHints: &storepb.QueryHints{ + ProjectionInclude: true, + ProjectionLabels: []string{"job", schema.SeriesHashColumn}, + }, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: true, + ProjectionLabels: []string{"job", schema.SeriesHashColumn}, + }, + }, + { + description: "honorProjectionHints=true, Mixed V1/V2 blocks, should reset projection", + honorProjectionHints: true, + queryHints: &storepb.QueryHints{ + ProjectionInclude: true, + ProjectionLabels: []string{"job"}, + }, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + createTestParquetBlock(t, false), // v1 + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: false, + ProjectionLabels: nil, + }, + }, + { + description: "honorProjectionHints=true, exclude projection (ProjectionInclude=false), should reset to full scan", + honorProjectionHints: true, + queryHints: &storepb.QueryHints{ + ProjectionInclude: false, + ProjectionLabels: []string{"job"}, + }, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: false, + ProjectionLabels: nil, // Non-include projections fall back to a full scan + }, + }, + { + description: "honorProjectionHints=true, nil query hints, should return default hints", + honorProjectionHints: true, + queryHints: nil, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: false, + ProjectionLabels: nil, + }, + }, + { + description: "honorProjectionHints=true, Empty projection labels, should add hash column only", + honorProjectionHints: true, + queryHints: &storepb.QueryHints{ + ProjectionInclude: true, + ProjectionLabels: []string{}, + }, + shards: []*parquetBlock{ + createTestParquetBlock(t, true), + }, + expectedHints: &prom_storage.SelectHints{ + Start: minT, + End: maxT, + ProjectionInclude: true, + ProjectionLabels: []string{schema.SeriesHashColumn}, + }, + }, + } + + for _, tc := range tests { + t.Run(tc.description, func(t *testing.T) { + store := &parquetBucketStore{ + honorProjectionHints: tc.honorProjectionHints, + } + shards := tc.shards + hints := store.buildSelectHints(tc.queryHints, shards, minT, maxT) + require.Equal(t, tc.expectedHints, hints) + }) + } +} + +func TestParquetBucketStore_buildSelectHints_DoesNotMutateRequestLabels(t *testing.T) { + // Spare capacity lets an unclipped append write into the request's backing array. + reqLabels := make([]string, 1, 2) + reqLabels[0] = "job" + queryHints := &storepb.QueryHints{ProjectionInclude: true, ProjectionLabels: reqLabels} + + store := &parquetBucketStore{honorProjectionHints: true} + hints := store.buildSelectHints(queryHints, []*parquetBlock{createTestParquetBlock(t, true)}, 1000, 2000) + + require.Equal(t, []string{"job", schema.SeriesHashColumn}, hints.ProjectionLabels) + require.Equal(t, []string{"job"}, queryHints.ProjectionLabels) + require.Equal(t, "", reqLabels[:2][1], "request backing array must not be written") +} + +func TestParquetBucketStore_Series_ProjectionHints(t *testing.T) { + const ( + userID = "user-1" + numSeries = 100 + numSamples = 10 + ) + + ctx := context.Background() + tmpDir := t.TempDir() + storageDir := filepath.Join(tmpDir, "storage") + dataDir := filepath.Join(tmpDir, "data") + + storageCfg := cortex_tsdb.BlocksStorageConfig{ + UsersScanner: users.UsersScannerConfig{ + Strategy: users.UserScanStrategyList, + UpdateInterval: time.Second, + }, + Bucket: bucket.Config{ + Backend: "filesystem", + Filesystem: filesystem.Config{Directory: storageDir}, + }, + BucketStore: cortex_tsdb.BucketStoreConfig{ + SyncDir: filepath.Join(tmpDir, "sync"), + BucketStoreType: "parquet", + BlockDiscoveryStrategy: string(cortex_tsdb.RecursiveDiscovery), + // Must be > 0, otherwise the per-query errgroup limit becomes 0 and + // errGroup.Go blocks forever, hanging the Series handler. + ParquetQueryConcurrency: 4, + }, + } + + bucketClient, err := bucket.NewClient(ctx, storageCfg.Bucket, nil, "test", log.NewNopLogger(), prometheus.NewRegistry()) + require.NoError(t, err) + + // Generated series carry these labels: __name__=test_metric, idx=, job=test_job, instance=localhost:9090. + blockID := prepareParquetBlock(t, ctx, storageCfg, bucketClient, dataDir, userID, numSeries, numSamples) + + startGRPCServer := func(honorProjectionHints bool) (storepb.StoreClient, func()) { + cfg := storageCfg + cfg.BucketStore.HonorProjectionHints = honorProjectionHints + + reg := prometheus.NewPedanticRegistry() + stores, err := NewBucketStores(cfg, NewNoShardingStrategy(log.NewNopLogger(), nil), objstore.WithNoopInstr(bucketClient), defaultLimitsOverrides(nil), mockLoggingLevel(), log.NewNopLogger(), reg) + require.NoError(t, err) + + listener, err := net.Listen("tcp", "localhost:0") + require.NoError(t, err) + + gRPCServer := grpc.NewServer(grpc.StreamInterceptor(middleware.StreamServerUserHeaderInterceptor)) + storepb.RegisterStoreServer(gRPCServer, stores) + go func() { + if err := gRPCServer.Serve(listener); err != nil && err != grpc.ErrServerStopped { + t.Error(err) + } + }() + + conn, err := grpc.NewClient(listener.Addr().String(), + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithDefaultCallOptions(grpc.MaxCallRecvMsgSize(math.MaxInt32)), + ) + require.NoError(t, err) + + return storepb.NewStoreClient(conn), func() { + _ = conn.Close() + gRPCServer.Stop() + } + } + + // querySeriesLabels issues a Series request and returns the union of all label names seen across + // the returned series. + querySeriesLabels := func(t *testing.T, client storepb.StoreClient, projectionLabels []string) map[string]struct{} { + reqCtx := grpcMetadata.NewOutgoingContext(context.Background(), grpcMetadata.Pairs(cortex_tsdb.TenantIDExternalLabel, userID)) + reqCtx, err := user.InjectIntoGRPCRequest(user.InjectOrgID(reqCtx, userID)) + require.NoError(t, err) + + seriesHints := &hintspb.SeriesRequestHints{ + BlockMatchers: []storepb.LabelMatcher{{Type: storepb.LabelMatcher_RE, Name: block.BlockIDLabel, Value: blockID}}, + } + hintsAny, err := types.MarshalAny(seriesHints) + require.NoError(t, err) + + req := &storepb.SeriesRequest{ + MinTime: 0, + MaxTime: math.MaxInt64, + Matchers: []storepb.LabelMatcher{{Type: storepb.LabelMatcher_RE, Name: "__name__", Value: ".+"}}, + ResponseBatchSize: 1000, + Hints: hintsAny, + } + if len(projectionLabels) > 0 { + req.QueryHints = &storepb.QueryHints{ProjectionInclude: true, ProjectionLabels: projectionLabels} + } + + stream, err := client.Series(reqCtx, req) + require.NoError(t, err) + + seen := map[string]struct{}{} + seriesCount := 0 + collect := func(zlabels []labelpb.ZLabel) { + seriesCount++ + for _, l := range zlabels { + seen[l.Name] = struct{}{} + } + } + for { + resp, err := stream.Recv() + if err == io.EOF { + break + } + require.NoError(t, err) + if s := resp.GetSeries(); s != nil { + collect(s.Labels) + } else if b := resp.GetBatch(); b != nil { + for _, s := range b.Series { + collect(s.Labels) + } + } + } + require.Equal(t, numSeries, seriesCount, "all series should be returned regardless of projection") + return seen + } + + client, stop := startGRPCServer(true) + defer stop() + + testCases := []struct { + name string + projectionLabels []string + wantPresent []string + wantAbsent []string + }{ + { + name: "without projection hints returns the full label set", + wantPresent: []string{"__name__", "job", "instance", "idx"}, + }, + { + name: "with projection hints only materializes requested labels", + projectionLabels: []string{"job"}, + wantPresent: []string{"job"}, + wantAbsent: []string{"__name__", "instance", "idx"}, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + labels := querySeriesLabels(t, client, tc.projectionLabels) + for _, l := range tc.wantPresent { + require.Contains(t, labels, l, "expected label %q to be materialized", l) + } + for _, l := range tc.wantAbsent { + require.NotContains(t, labels, l, "expected label %q to be projected away", l) + } + }) + } +} diff --git a/pkg/storegateway/parquet_bucket_stores.go b/pkg/storegateway/parquet_bucket_stores.go index 9c6ac89bec5..1ab50dfaa40 100644 --- a/pkg/storegateway/parquet_bucket_stores.go +++ b/pkg/storegateway/parquet_bucket_stores.go @@ -298,17 +298,18 @@ func (u *ParquetBucketStores) createParquetBucketStore(userID string, userLogger userBucket := bucket.NewUserBucketClient(userID, u.bucket, u.limits) store := &parquetBucketStore{ - logger: userLogger, - bucket: userBucket, - indexLoader: u.indexLoader, - limits: u.limits, - userID: userID, - bucketIndexEnabled: u.cfg.BucketStore.BucketIndex.Enabled, - concurrency: u.cfg.BucketStore.ParquetQueryConcurrency, - chunksDecoder: u.chunksDecoder, - matcherCache: u.matcherCache, - parquetShardCache: u.parquetShardCache, - rowRangesCache: u.rowRangesCache, + logger: userLogger, + bucket: userBucket, + indexLoader: u.indexLoader, + limits: u.limits, + userID: userID, + bucketIndexEnabled: u.cfg.BucketStore.BucketIndex.Enabled, + concurrency: u.cfg.BucketStore.ParquetQueryConcurrency, + chunksDecoder: u.chunksDecoder, + matcherCache: u.matcherCache, + parquetShardCache: u.parquetShardCache, + honorProjectionHints: u.cfg.BucketStore.HonorProjectionHints, + rowRangesCache: u.rowRangesCache, } return store, nil @@ -358,7 +359,7 @@ func (p *parquetBucketStore) newParquetBlock(ctx context.Context, name string, s if err != nil { return nil, err } - m, err := search.NewMaterializer(s, d, shard, p.concurrency, rowCountQuota, chunkBytesQuota, dataBytesQuota, search.NoopMaterializedSeriesFunc, materializedLabelsFilterCallback, false) + m, err := search.NewMaterializer(s, d, shard, p.concurrency, rowCountQuota, chunkBytesQuota, dataBytesQuota, search.NoopMaterializedSeriesFunc, materializedLabelsFilterCallback, p.honorProjectionHints) if err != nil { return nil, err } @@ -414,7 +415,7 @@ func (f *shardMatcherLabelsFilter) Close() { f.shardMatcher.Close() } -func (b *parquetBlock) Query(ctx context.Context, mint, maxt int64, skipChunks bool, matchers []*labels.Matcher) (prom_storage.ChunkSeriesSet, error) { +func (b *parquetBlock) Query(ctx context.Context, hints *prom_storage.SelectHints, skipChunks bool, matchers []*labels.Matcher) (prom_storage.ChunkSeriesSet, error) { errGroup, ctx := errgroup.WithContext(ctx) errGroup.SetLimit(b.concurrency) @@ -444,7 +445,7 @@ func (b *parquetBlock) Query(ctx context.Context, mint, maxt int64, skipChunks b return nil } - seriesSetIter, err := b.m.Materialize(ctx, nil, rgi, mint, maxt, skipChunks, rr) + seriesSetIter, err := b.m.Materialize(ctx, hints, rgi, hints.Start, hints.End, skipChunks, rr) if err != nil { return err } @@ -581,6 +582,18 @@ func (b *parquetBlock) allLabelValues(ctx context.Context, name string, limit in return util.MergeUnsortedSlices(int(limit), results...), nil } +// hasHashColumn checks if the parquet block contains the schema.SeriesHashColumn column. +// This is used to determine if projection pushdown can be enabled. +func (b *parquetBlock) hasHashColumn() bool { + labelsFile := b.shard.LabelsFile() + if labelsFile == nil { + return false + } + + _, found := labelsFile.Schema().Lookup(schema.SeriesHashColumn) + return found +} + type byLabels []prom_storage.ChunkSeries func (b byLabels) Len() int { return len(b) } diff --git a/schemas/cortex-config-schema.json b/schemas/cortex-config-schema.json index 70c01e2a1c5..35d8c616a25 100644 --- a/schemas/cortex-config-schema.json +++ b/schemas/cortex-config-schema.json @@ -1539,6 +1539,12 @@ "x-cli-flag": "blocks-storage.bucket-store.consistency-delay", "x-format": "duration" }, + "honor_projection_hints": { + "default": false, + "description": "[Experimental] If enabled, Store Gateway will honor projection hints and only materialize requested labels. It only takes effect when `-blocks-storage.bucket-store.bucket-store-type` is parquet and `-querier.honor-projection-hints` is enabled.", + "type": "boolean", + "x-cli-flag": "blocks-storage.bucket-store.honor-projection-hints" + }, "ignore_blocks_before": { "default": "0s", "description": "The blocks created before `now() - ignore_blocks_before` will not be synced. 0 to disable.", From f10a22f64ccd663459bab20192c84a46cf95755e Mon Sep 17 00:00:00 2001 From: SungJin1212 Date: Tue, 6 Oct 2026 20:39:39 +0900 Subject: [PATCH 2/3] update doc Signed-off-by: SungJin1212 --- docs/blocks-storage/querier.md | 6 ++++-- docs/blocks-storage/store-gateway.md | 6 ++++-- docs/configuration/config-file-reference.md | 6 ++++-- pkg/storage/tsdb/config.go | 2 +- schemas/cortex-config-schema.json | 2 +- 5 files changed, 14 insertions(+), 8 deletions(-) diff --git a/docs/blocks-storage/querier.md b/docs/blocks-storage/querier.md index 5796d036b73..402718c128d 100644 --- a/docs/blocks-storage/querier.md +++ b/docs/blocks-storage/querier.md @@ -2168,8 +2168,10 @@ blocks_storage: # [Experimental] If enabled, Store Gateway will honor projection hints and # only materialize requested labels. It only takes effect when - # `-blocks-storage.bucket-store.bucket-store-type` is parquet and - # `-querier.honor-projection-hints` is enabled. + # `-blocks-storage.bucket-store.bucket-store-type` is parquet, + # `-querier.honor-projection-hints` is enabled and the querier uses the + # Thanos engine (`-querier.thanos-engine`) with the `projection` optimizer + # (`-querier.optimizers`). # CLI flag: -blocks-storage.bucket-store.honor-projection-hints [honor_projection_hints: | default = false] diff --git a/docs/blocks-storage/store-gateway.md b/docs/blocks-storage/store-gateway.md index 4ff4ef2836f..a8cafa6c26d 100644 --- a/docs/blocks-storage/store-gateway.md +++ b/docs/blocks-storage/store-gateway.md @@ -2221,8 +2221,10 @@ blocks_storage: # [Experimental] If enabled, Store Gateway will honor projection hints and # only materialize requested labels. It only takes effect when - # `-blocks-storage.bucket-store.bucket-store-type` is parquet and - # `-querier.honor-projection-hints` is enabled. + # `-blocks-storage.bucket-store.bucket-store-type` is parquet, + # `-querier.honor-projection-hints` is enabled and the querier uses the + # Thanos engine (`-querier.thanos-engine`) with the `projection` optimizer + # (`-querier.optimizers`). # CLI flag: -blocks-storage.bucket-store.honor-projection-hints [honor_projection_hints: | default = false] diff --git a/docs/configuration/config-file-reference.md b/docs/configuration/config-file-reference.md index 9f69f5d6201..938859d3af2 100644 --- a/docs/configuration/config-file-reference.md +++ b/docs/configuration/config-file-reference.md @@ -2865,8 +2865,10 @@ bucket_store: # [Experimental] If enabled, Store Gateway will honor projection hints and # only materialize requested labels. It only takes effect when - # `-blocks-storage.bucket-store.bucket-store-type` is parquet and - # `-querier.honor-projection-hints` is enabled. + # `-blocks-storage.bucket-store.bucket-store-type` is parquet, + # `-querier.honor-projection-hints` is enabled and the querier uses the Thanos + # engine (`-querier.thanos-engine`) with the `projection` optimizer + # (`-querier.optimizers`). # CLI flag: -blocks-storage.bucket-store.honor-projection-hints [honor_projection_hints: | default = false] diff --git a/pkg/storage/tsdb/config.go b/pkg/storage/tsdb/config.go index 7bbcddcc66c..e173552cebb 100644 --- a/pkg/storage/tsdb/config.go +++ b/pkg/storage/tsdb/config.go @@ -406,7 +406,7 @@ func (cfg *BucketStoreConfig) RegisterFlags(f *flag.FlagSet) { f.IntVar(&cfg.MatchersCacheMaxItems, "blocks-storage.bucket-store.matchers-cache-max-items", 0, "Maximum number of entries in the regex matchers cache. 0 to disable.") f.IntVar(&cfg.ParquetQueryConcurrency, "blocks-storage.bucket-store.parquet-query-concurrency", 4, "Maximum number of concurrent goroutines per query applied at each level of parquet processing: shard querying, row group processing, and column materialization. Note: this limit is applied independently at each level, so the total goroutines per query can grow multiplicatively (up to N^3 in the worst case).") cfg.ParquetShardCache.RegisterFlagsWithPrefix("blocks-storage.bucket-store.", f) - f.BoolVar(&cfg.HonorProjectionHints, "blocks-storage.bucket-store.honor-projection-hints", false, "[Experimental] If enabled, Store Gateway will honor projection hints and only materialize requested labels. It only takes effect when `-blocks-storage.bucket-store.bucket-store-type` is parquet and `-querier.honor-projection-hints` is enabled.") + f.BoolVar(&cfg.HonorProjectionHints, "blocks-storage.bucket-store.honor-projection-hints", false, "[Experimental] If enabled, Store Gateway will honor projection hints and only materialize requested labels. It only takes effect when `-blocks-storage.bucket-store.bucket-store-type` is parquet, `-querier.honor-projection-hints` is enabled and the querier uses the Thanos engine (`-querier.thanos-engine`) with the `projection` optimizer (`-querier.optimizers`).") } // Validate the config. diff --git a/schemas/cortex-config-schema.json b/schemas/cortex-config-schema.json index 35d8c616a25..c9361f8e27e 100644 --- a/schemas/cortex-config-schema.json +++ b/schemas/cortex-config-schema.json @@ -1541,7 +1541,7 @@ }, "honor_projection_hints": { "default": false, - "description": "[Experimental] If enabled, Store Gateway will honor projection hints and only materialize requested labels. It only takes effect when `-blocks-storage.bucket-store.bucket-store-type` is parquet and `-querier.honor-projection-hints` is enabled.", + "description": "[Experimental] If enabled, Store Gateway will honor projection hints and only materialize requested labels. It only takes effect when `-blocks-storage.bucket-store.bucket-store-type` is parquet, `-querier.honor-projection-hints` is enabled and the querier uses the Thanos engine (`-querier.thanos-engine`) with the `projection` optimizer (`-querier.optimizers`).", "type": "boolean", "x-cli-flag": "blocks-storage.bucket-store.honor-projection-hints" }, From 6792a7b43b815221cee312246c99997f9d9a95da Mon Sep 17 00:00:00 2001 From: SungJin1212 Date: Tue, 6 Oct 2026 20:49:58 +0900 Subject: [PATCH 3/3] fix lint Signed-off-by: SungJin1212 --- integration/parquet_store_gateway_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/integration/parquet_store_gateway_test.go b/integration/parquet_store_gateway_test.go index 455d138a5cd..5f33ae77122 100644 --- a/integration/parquet_store_gateway_test.go +++ b/integration/parquet_store_gateway_test.go @@ -110,7 +110,7 @@ func TestParquetBucketStore_ProjectionHint(t *testing.T) { numSamples := 100 lbls := make([]labels.Labels, 0, numSeries) - for i := 0; i < numSeries; i++ { + for i := range numSeries { lbls = append(lbls, labels.FromStrings( labels.MetricName, "http_requests_total", "job", "api-server", @@ -145,7 +145,7 @@ func TestParquetBucketStore_ProjectionHint(t *testing.T) { cHonorOff, err := e2ecortex.NewClient("", querierHonorOff.HTTPEndpoint(), "", "", "user-1") require.NoError(t, err) - cortex_testutil.Poll(t, 60*time.Second, true, func() interface{} { + cortex_testutil.Poll(t, 60*time.Second, true, func() any { labelSets, err := c.Series([]string{`{job="api-server"}`}, start, end) if err != nil { t.Logf("Series query failed: %v", err)