Skip to content

Fix #214: skip end-offset metric when highwater is unknown (None)#703

Open
wbarnha wants to merge 1 commit into
masterfrom
claude/fix-214-highwater-none-metric
Open

Fix #214: skip end-offset metric when highwater is unknown (None)#703
wbarnha wants to merge 1 commit into
masterfrom
claude/fix-214-highwater-none-metric

Conversation

@wbarnha

@wbarnha wbarnha commented Jul 19, 2026

Copy link
Copy Markdown
Member

What

During a rebalance the consumer could crash _drain_messages with:

TypeError: float() argument must be a string or a number, not 'NoneType'
  File ".../faust/transport/consumer.py", line 715, in getmany
    self.app.monitor.track_tp_end_offset(tp, highwater_mark)
  File ".../faust/sensors/prometheus.py", line 496, in track_tp_end_offset
    self._metrics.topic_partition_end_offset.labels(...).set(float(value))

Fixes #214.

Why

Consumer.highwater(tp) legitimately returns None during a rebalance, before the partition's end offset is known. getmany() passed that value straight into monitor.track_tp_end_offset(tp, highwater_mark). The base Monitor stores it harmlessly, but the PrometheusMonitor / DatadogMonitor / StatsdMonitor implementations call float(offset), so a None raises TypeError and takes down the drain loop (the AeroSpike/rebalance framing in the report is incidental — any storage backend can hit this).

How

Guard the call in getmany() so the end-offset metric is only tracked once the highwater is known. Fixing it at the single source point covers every metric sensor at once, rather than patching float() in each sensor separately, and is semantically correct — there is no meaningful end offset to record yet.

Test

Added two tests in tests/unit/transport/test_consumer.py:

  • test_getmany__highwater_none_not_trackedhighwater() returns None; asserts messages still flow and track_tp_end_offset is not called (verified to fail before the fix).
  • test_getmany__highwater_trackedhighwater() returns an int; asserts track_tp_end_offset is called with it.

Full tests/unit/transport/test_consumer.py passes (115 passed, 2 skipped).

🤖 Generated with Claude Code


Generated by Claude Code

During a rebalance Consumer.highwater(tp) can return None before the end
offset is known.  getmany() passed that straight to
monitor.track_tp_end_offset(), and the Prometheus/Datadog/StatsD sensors
do float(offset) on it, raising
'TypeError: float() argument must be a string or a number, not NoneType'
inside _drain_messages and crashing the consumer.

Guard the call so the end-offset metric is only tracked once the highwater
is known.  This fixes every metric sensor at the single source instead of
patching each one.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL
@codecov

codecov Bot commented Jul 19, 2026

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 94.15%. Comparing base (3073eb9) to head (e6871b2).
⚠️ Report is 3 commits behind head on master.

Additional details and impacted files
@@           Coverage Diff           @@
##           master     #703   +/-   ##
=======================================
  Coverage   94.14%   94.15%           
=======================================
  Files         104      104           
  Lines       11136    11137    +1     
  Branches     1201     1202    +1     
=======================================
+ Hits        10484    10486    +2     
+ Misses        551      550    -1     
  Partials      101      101           

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

wbarnha added a commit that referenced this pull request Jul 19, 2026
Per review, the v0.12.0 changelog/release notes should describe only what is
already on master, not work still in open PRs.

- Remove the not-yet-merged items: the offset-commit data-loss fixes
  (#606/#707, #316/#692), the optional OpenTracing/OpenTelemetry extras
  (#685/#686, #688/#681), web_application_options (#704), and the reported-issue
  fix stack (#693-#703, #705). These will be added back as they merge.
- Add a Dependencies section noting the current runtime/client libraries:
  mode-streaming >= 0.4.0, aiokafka >= 0.10.0 (compatible with recent 0.13/0.14
  releases), the new confluent-kafka >= 2.0.0 for faust[ckafka], and the
  faust-cchardet fork replacing unmaintained cchardet.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL
wbarnha added a commit to SpencerWhitehead7/faust that referenced this pull request Jul 21, 2026
…st-streaming#708)

* docs: prepare v0.12.0 release notes and changelog

Resume the Keep a Changelog format (dormant since v0.8.10) with a v0.12.0
section, and add standalone GitHub release notes covering the changes since
v0.11.3 plus the pending fix stack.

Highlights: two offset data-loss fixes (faust-streaming#606/faust-streaming#707, faust-streaming#316/faust-streaming#692), the
re-added confluent-kafka driver, Python 3.14 support (3.8/3.9 dropped),
OpenTracing/OpenTelemetry made optional, and a live-broker CI harness.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL

* docs: drop closed codecov.yml PR (faust-streaming#683) from v0.12.0 notes

PR faust-streaming#683 (codecov.yml with a 1% coverage threshold) was closed without
merging, so remove it from the changelog and release notes to keep the
v0.12.0 change list accurate.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL

* docs: derive Sphinx version from the package instead of a stale constant

`docs/conf.py` hardcoded `version_dev='1.1'` / `version_stable='1.0'` -
robinhood-era values that never matched faust-streaming's 0.x line, so the
published GitHub Pages docs advertised the wrong version.

Derive the documented major.minor from `faust.__version__` (which
setuptools_scm resolves from the git tag), so the docs always report the real
version and this can't silently drift between releases.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL

* docs: scope v0.12.0 notes to merged work; note dependency updates

Per review, the v0.12.0 changelog/release notes should describe only what is
already on master, not work still in open PRs.

- Remove the not-yet-merged items: the offset-commit data-loss fixes
  (faust-streaming#606/faust-streaming#707, faust-streaming#316/faust-streaming#692), the optional OpenTracing/OpenTelemetry extras
  (faust-streaming#685/faust-streaming#686, faust-streaming#688/faust-streaming#681), web_application_options (faust-streaming#704), and the reported-issue
  fix stack (faust-streaming#693-faust-streaming#703, faust-streaming#705). These will be added back as they merge.
- Add a Dependencies section noting the current runtime/client libraries:
  mode-streaming >= 0.4.0, aiokafka >= 0.10.0 (compatible with recent 0.13/0.14
  releases), the new confluent-kafka >= 2.0.0 for faust[ckafka], and the
  faust-cchardet fork replacing unmaintained cchardet.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL

* docs: set v0.12.0 changelog date to 2026-07-19

Replace the UNRELEASED placeholder with the release date and drop the
now-satisfied "set the date at tag time" note.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HHPL4VFWQRQPpjR1gXSKyL

* Delete RELEASE_NOTES_v0.12.0.md

---------

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Faust *sometimes* crashes during rebalance when AeroSpike is used as storage engine

1 participant