apache / apache/iggy

bug(connectors): PollingMessages auto-commit commits offsets before sink processing

Open
#2,928 3 comments 0 reactions 0 assignees View on GitHub
bug connectors rust
Dominant language
Rust
Stars
4.9k
Forks
432
Avg merge
2d 10h
Merged PRs (30d)
173

Description

## Description

The `AutoCommitWhen::PollingMessages` strategy commits consumer group offsets **before** `consume()` is called on the sink plugin. This means if `consume()` fails (and the failure is already silently discarded per #2927), the messages are permanently lost — the consumer group has already advanced past them.

## Sequence

```
1. consumer.next() → messages received
2. Offsets committed → consumer group advances ← PROBLEM
3. process_messages() → calls sink consume() via FFI
4. consume() fails → messages already committed, lost forever
```

## Location

`core/connectors/runtime/src/sink.rs`:

- **Line 421**: Consumer configured with `AutoCommitWhen::PollingMessages`
```rust
.auto_commit(AutoCommit::When(AutoCommitWhen::PollingMessages))
```

- **Lines 266-344** (`consume_messages` loop): `consumer.next()` at line 272 triggers the auto-commit, but `process_messages()` doesn't execute until line 311.

## Impact

Combined with #2927 (consume return value discarded), **at-least-once delivery is not achievable** with the current runtime for any sink connector. If `consume()` fails and the process restarts, those messages will never be retried because the offsets have already been committed.

## Suggested Fix

Change auto-commit strategy to `AutoCommitWhen::ConsumerStopped` or `AutoCommit::Disabled`, and commit offsets **after** successful `consume()`:

```rust
// Option A: Commit after processing
.auto_commit(AutoCommit::Disabled)
// ... in the consume loop, after successful consume():
consumer.store_offset(offset).await?;

// Option B: Use AfterPollingMessages if available
.auto_commit(AutoCommit::When(AutoCommitWhen::AfterProcessing))
```

The exact API depends on the Iggy SDK's consumer group offset management capabilities.

## Related

- #2927 (consume return value discarded)
- Discussed in #2919 (HTTP sink connector proposal)
- HTTP sink PR: #2925

Contributor guide

Open the contributing guide

Research direction

Start in core/connectors/runtime/src/sink.rs, reading the consume_messages loop at lines 266-344 and the consumer configuration at line 421. Check the Iggy SDK offset-management API and related issues #2927 and #2919 before choosing the commit point. Done means successful sink processing is followed by an offset commit, while failed processing leaves messages eligible for retry.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
backend, distributed-systems
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
48/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.