Skip to content

Commit 4b9a6c5

Browse files
committed
Implement auto-forget in Parquet converter lifecycle (cortexproject#7752)
* implement auto-forget in parquet converter ring lifecycler Signed-off-by: SungJin1212 <tjdwls1201@gmail.com> * clean annotations Signed-off-by: SungJin1212 <tjdwls1201@gmail.com> --------- Signed-off-by: SungJin1212 <tjdwls1201@gmail.com>
1 parent 4edf2f7 commit 4b9a6c5

3 files changed

Lines changed: 59 additions & 1 deletion

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,7 @@
8787
* [BUGFIX] Querier: Fix gRPC `codes.Canceled` errors being mapped to HTTP 500 instead of 499 when a client cancels a query. #7738
8888
* [BUGFIX] Compactor: Fix spurious `bucket operation fail after retries` error logs emitted during partial block cleanup. #7749
8989
* [BUGFIX] Alertmanager: Fix panic in `validateAlertmanagerConfig` when receiver config traversal encounters nil interface values. #7751
90+
* [BUGFIX] Parquet Converter: Fix `auto_forget_delay` having no effect. The ring lifecycler was created without the auto-forget delegate, so unhealthy instances were never automatically removed from the ring. #7752
9091

9192
## 1.21.1 2026-06-04
9293

pkg/parquetconverter/converter.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -174,7 +174,12 @@ func newConverter(cfg Config, bkt objstore.InstrumentedBucket, storageCfg cortex
174174
func (c *Converter) starting(ctx context.Context) error {
175175
lifecyclerCfg := c.cfg.Ring.ToLifecyclerConfig()
176176
var err error
177-
c.ringLifecycler, err = ring.NewLifecycler(lifecyclerCfg, ring.NewNoopFlushTransferer(), "parquet-converter", ringKey, true, false, c.logger, prometheus.WrapRegistererWithPrefix("cortex_", c.reg))
177+
var delegate ring.LifecyclerDelegate
178+
delegate = &ring.DefaultLifecyclerDelegate{}
179+
if c.cfg.Ring.AutoForgetDelay > 0 {
180+
delegate = ring.NewLifecyclerAutoForgetDelegate(c.cfg.Ring.AutoForgetDelay, delegate, c.logger)
181+
}
182+
c.ringLifecycler, err = ring.NewLifecyclerWithDelegate(lifecyclerCfg, ring.NewNoopFlushTransferer(), "parquet-converter", ringKey, true, false, c.logger, prometheus.WrapRegistererWithPrefix("cortex_", c.reg), delegate)
178183
if err != nil {
179184
return errors.Wrap(err, "unable to initialize converter ring lifecycler")
180185
}

pkg/parquetconverter/converter_test.go

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -682,3 +682,55 @@ func TestConvertWithMaxNumColumns(t *testing.T) {
682682
require.NoError(t, err)
683683
require.Equal(t, 1, shards2, "expected single shard with high column limit")
684684
}
685+
686+
func TestConverter_RingLifecyclerShouldAutoForgetUnhealthyInstances(t *testing.T) {
687+
// Create a shared KV Store
688+
kvstore, closer := consul.NewInMemoryClient(ring.GetCodec(), log.NewNopLogger(), nil)
689+
t.Cleanup(func() { assert.NoError(t, closer.Close()) })
690+
691+
bucketClient := objstore.WithNoopInstr(objstore.NewInMemBucket())
692+
limits := &validation.Limits{}
693+
flagext.DefaultValues(limits)
694+
limits.ParquetConverterEnabled = true
695+
696+
// Create two converters
697+
var converters []*Converter
698+
for i := range 2 {
699+
cfg := prepareConfig()
700+
cfg.Ring.InstanceID = fmt.Sprintf("parquet-converter-%d", i)
701+
cfg.Ring.InstanceAddr = fmt.Sprintf("127.0.0.%d", i+1)
702+
cfg.Ring.KVStore.Mock = kvstore
703+
cfg.Ring.HeartbeatPeriod = 100 * time.Millisecond
704+
cfg.Ring.HeartbeatTimeout = 200 * time.Millisecond
705+
cfg.Ring.AutoForgetDelay = 400 * time.Millisecond
706+
707+
c, _, _ := prepare(t, cfg, bucketClient, limits, nil)
708+
converters = append(converters, c)
709+
}
710+
711+
// Start both converters.
712+
require.NoError(t, services.StartAndAwaitRunning(context.Background(), converters[0]))
713+
require.NoError(t, services.StartAndAwaitRunning(context.Background(), converters[1]))
714+
715+
// Both should be healthy.
716+
test.Poll(t, 5*time.Second, true, func() any {
717+
healthy, unhealthy, _ := converters[0].ring.GetAllInstanceDescs(ring.Reporting)
718+
return len(healthy) == 2 && len(unhealthy) == 0
719+
})
720+
721+
// Override UnregisterOnShutdown so the instance stays in the ring after stop,
722+
// simulating a crash or ungraceful shutdown.
723+
converters[1].ringLifecycler.SetUnregisterOnShutdown(false)
724+
// The converter running() returns ctx.Err() on stop, so context.Canceled is expected.
725+
err := services.StopAndAwaitTerminated(context.Background(), converters[1])
726+
require.True(t, err == nil || errors.Is(err, context.Canceled), "unexpected error stopping converter: %v", err)
727+
728+
// The stopped instance should appear unhealthy first, then be auto-forgotten.
729+
test.Poll(t, 5*time.Second, true, func() any {
730+
healthy, unhealthy, _ := converters[0].ring.GetAllInstanceDescs(ring.Reporting)
731+
return len(healthy) == 1 && len(unhealthy) == 0
732+
})
733+
734+
err = services.StopAndAwaitTerminated(context.Background(), converters[0])
735+
require.True(t, err == nil || errors.Is(err, context.Canceled), "unexpected error stopping converter: %v", err)
736+
}

0 commit comments

Comments
 (0)