typelevel / typelevel/fs2

Feature request: Dynamic metering of Streams

Open
#3,328 3 comments 1 reaction 0 assignees View on GitHub

Nobody has claimed this yet.

Dominant language
Scala
Stars
2.5k
Forks
636
Avg merge
2d 4h
Merged PRs (30d)
7

Description

I initialy raised this in fs2-kafka, but I think this actually belongs here;

https://github.com/fd4s/fs2-kafka/issues/1270

In the same manner that you can use a Signal[F, Boolean] to pause consumption, which I've found incredibly useful for code with e.g. dynamic feature flags to turn on/off consumption, I am hoping to throttle consumption.

If this would be useful to other people, would you be open to a PR for this? Or, if this already exists, please do point me in its direction 🙏

1. is there currently a baked-in way to have dynamic metering on a stream, other than `evalTap(_ => doSomeDelay())

The limit of this is that someSleep won't have access to when the stream last emitted, how much time has elapsed, etc etc, which leads me to my actual question:

2. I think this is the feature request I'm asking for, if people would find this useful

// ignoring the details of where these signals come from
val targetDelay: Signal[F, FiniteDuration] = createDelaySignal[F]() // edit: see below, signal mightn't be the best choice here
val pauseSignal: Signal[F, Boolean] = createPauseSignal[F]()

// I would like to be able to:
someStream.
    .pauseWhen(pauseSignal)
    .meteredBy(targetDelay)
   .flatMap(...etcetc)

I would want the semantics of it to internally keep track of last time it produced a message, and if signal value is emitted that is less than what was last set, AND more time has passed since it was set, then it would immediately emit and then wait for the next duration to elapse. And inversely, if the duration is increased since last, then the time-delta is taken into consideration

So - in BDD format - this is the behaviour that I'm after:

Scenario 1: Initial startup
GIVEN target delay is initially 1 minute
WHEN the stream is started
THEN it emits a message every minute as long as there are more elements to emit

Scenario 2: Delay is changed to be < than it was previously
GIVEN targetDelay is initially 1 minute
AND the stream emits one message
AND 30 seconds pass
WHEN the targetDelay signal is changed to 5 seconds
THEN a message is immediately emitted because (newTime - timePassedSinceLastInvocaton) is negative or 0
AND every 5 seconds after this, a new message is emitted

Scenario 3: Delay is changed to be > than what it was previously
GIVEN targetDelay is initially 1 minute
AND the stream emits one message
AND 30 seconds pass
WHEN the targetDelay signal is changed to 45 seconds
THEN after an additional 15 seconds elapse, a message is emitted
AND every 45 seconds after this, a new message is emitted

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start by reading the linked fs2-kafka issue and the existing Signal-based pauseWhen behavior. Compare the proposed meteredBy(targetDelay) API with the three BDD scenarios, including shortened and lengthened delays. Done means the project has agreed on the API and semantics and the documented scenarios are covered by implementation tests.

Written by the indexing model from the issue text.

Assessment

Tech stack
scala
Domain
stream-processing
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.