(t *testing.T, cfg querierShardingTestConfig)
| 53 | } |
| 54 | |
| 55 | func runQuerierShardingTest(t *testing.T, cfg querierShardingTestConfig) { |
| 56 | // Going to high starts hitting file descriptor limit, since we run all queriers concurrently. |
| 57 | const batchSize = 100 |
| 58 | const numQueries = 500 |
| 59 | |
| 60 | s, err := e2e.NewScenario(networkName) |
| 61 | require.NoError(t, err) |
| 62 | defer s.Close() |
| 63 | |
| 64 | memcached := e2ecache.NewMemcached() |
| 65 | consul := e2edb.NewConsul() |
| 66 | require.NoError(t, s.StartAndWaitReady(consul, memcached)) |
| 67 | |
| 68 | flags := mergeFlags(BlocksStorageFlags(), map[string]string{ |
| 69 | "-querier.cache-results": "true", |
| 70 | "-querier.split-queries-by-interval": "24h", |
| 71 | "-querier.query-ingesters-within": "12h", // Required by the test on query /series out of ingesters time range |
| 72 | "-frontend.memcached.addresses": "dns+" + memcached.NetworkEndpoint(e2ecache.MemcachedPort), |
| 73 | "-frontend.max-outstanding-requests-per-tenant": strconv.Itoa(numQueries), // To avoid getting errors. |
| 74 | }) |
| 75 | |
| 76 | minio := e2edb.NewMinio(9000, flags["-blocks-storage.s3.bucket-name"]) |
| 77 | require.NoError(t, s.StartAndWaitReady(minio)) |
| 78 | |
| 79 | if cfg.shuffleShardingEnabled { |
| 80 | // Use only single querier for each user. |
| 81 | flags["-frontend.max-queriers-per-tenant"] = "1" |
| 82 | } |
| 83 | |
| 84 | // Start the query-scheduler if enabled. |
| 85 | var queryScheduler *e2ecortex.CortexService |
| 86 | if cfg.querySchedulerEnabled { |
| 87 | queryScheduler = e2ecortex.NewQueryScheduler("query-scheduler", flags, "") |
| 88 | require.NoError(t, s.StartAndWaitReady(queryScheduler)) |
| 89 | flags["-frontend.scheduler-address"] = queryScheduler.NetworkGRPCEndpoint() |
| 90 | flags["-querier.scheduler-address"] = queryScheduler.NetworkGRPCEndpoint() |
| 91 | } |
| 92 | |
| 93 | // Start the query-frontend. |
| 94 | queryFrontend := e2ecortex.NewQueryFrontend("query-frontend", flags, "") |
| 95 | require.NoError(t, s.Start(queryFrontend)) |
| 96 | |
| 97 | if !cfg.querySchedulerEnabled { |
| 98 | flags["-querier.frontend-address"] = queryFrontend.NetworkGRPCEndpoint() |
| 99 | } |
| 100 | |
| 101 | // Start all other services. |
| 102 | ingester := e2ecortex.NewIngester("ingester", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), flags, "") |
| 103 | distributor := e2ecortex.NewDistributor("distributor", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), flags, "") |
| 104 | querier1 := e2ecortex.NewQuerier("querier-1", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), flags, "") |
| 105 | querier2 := e2ecortex.NewQuerier("querier-2", e2ecortex.RingStoreConsul, consul.NetworkHTTPEndpoint(), flags, "") |
| 106 | |
| 107 | require.NoError(t, s.StartAndWaitReady(querier1, querier2, ingester, distributor)) |
| 108 | require.NoError(t, s.WaitReady(queryFrontend)) |
| 109 | |
| 110 | // Wait until distributor and queriers have updated the ring. |
| 111 | require.NoError(t, distributor.WaitSumMetrics(e2e.Equals(512), "cortex_ring_tokens_total")) |
| 112 | require.NoError(t, querier1.WaitSumMetrics(e2e.Equals(512), "cortex_ring_tokens_total")) |
no test coverage detected