-
Notifications
You must be signed in to change notification settings - Fork 498
[FLINK-33525] Move ImpulseSource to new Source API #950
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 4 commits
618d50a
de5fb01
0d2bd2d
eb8ef85
de713b7
fabd090
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -18,12 +18,16 @@ | |
|
|
||
| package autoscaling; | ||
|
|
||
| import org.apache.flink.api.common.eventtime.WatermarkStrategy; | ||
| import org.apache.flink.api.common.functions.RichFlatMapFunction; | ||
| import org.apache.flink.api.common.typeinfo.Types; | ||
| import org.apache.flink.api.connector.source.util.ratelimit.RateLimiterStrategy; | ||
| import org.apache.flink.api.java.utils.ParameterTool; | ||
| import org.apache.flink.connector.datagen.source.DataGeneratorSource; | ||
| import org.apache.flink.connector.datagen.source.GeneratorFunction; | ||
| import org.apache.flink.streaming.api.datastream.DataStream; | ||
| import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; | ||
| import org.apache.flink.streaming.api.functions.sink.DiscardingSink; | ||
| import org.apache.flink.streaming.api.functions.source.SourceFunction; | ||
| import org.apache.flink.util.Collector; | ||
|
|
||
| import org.slf4j.Logger; | ||
|
|
@@ -74,8 +78,50 @@ public static void main(String[] args) throws Exception { | |
| for (String branch : maxLoadPerTask.split("\n")) { | ||
| String[] taskLoads = branch.split(";"); | ||
|
|
||
| /* | ||
| * Creates an unbounded stream that continuously emits the constant value 42L. | ||
| * Flink's DataGeneratorSource with RateLimiterStrategy is used to control the emission rate. | ||
| * | ||
| * Emission Rate Logic: | ||
| * - The goal is to generate a fixed number of impulses per sampling interval. | ||
| * - `samplingIntervalMs` defines the duration of one sampling interval in milliseconds. | ||
| * - We define `IMPULSES_PER_SAMPLING_INTERVAL = 10`, meaning that for every sampling interval, | ||
| * exactly 10 impulses should be generated. | ||
| * | ||
| * To calculate the total number of records emitted per second: | ||
| * 1. Determine how many sampling intervals fit within one second: | ||
| * samplingIntervalsPerSecond = 1000 / samplingIntervalMs; | ||
| * 2. Multiply this by the number of impulses per interval to get the total rate: | ||
| * impulsesPerSecond = IMPULSES_PER_SAMPLING_INTERVAL * samplingIntervalsPerSecond; | ||
| * | ||
| * Example Calculations: | ||
| * - If `samplingIntervalMs = 1000 ms`: | ||
| * - `samplingIntervalsPerSecond = 1000 / 1000 = 1` | ||
| * - `impulsesPerSecond = 10 * 1 = 10 records per second` | ||
| * - If `samplingIntervalMs = 500 ms`: | ||
| * - `samplingIntervalsPerSecond = 1000 / 500 = 2` | ||
| * - `impulsesPerSecond = 10 * 2 = 20 records per second` | ||
| * - If `samplingIntervalMs = 2000 ms`: | ||
| * - `samplingIntervalsPerSecond = 1000 / 2000 = 0.5` | ||
| * - `impulsesPerSecond = 10 * 0.5 = 5 records per second` | ||
| * | ||
| * This approach ensures that the number of records emitted dynamically scales | ||
|
||
| * based on the sampling interval while maintaining the target of 10 impulses per interval. | ||
| * RateLimiterStrategy internally distributes these emissions efficiently over time. | ||
| */ | ||
| DataStream<Long> stream = | ||
| env.addSource(new ImpulseSource(samplingIntervalMs)).name("ImpulseSource"); | ||
| env.fromSource( | ||
| new DataGeneratorSource<>( | ||
| (GeneratorFunction<Long, Long>) | ||
| (index) -> 42L, // Emits constant value 42 | ||
| Long.MAX_VALUE, // Unbounded stream | ||
| RateLimiterStrategy.perSecond( | ||
| (double) 1000 | ||
|
||
| / ((double) samplingIntervalMs | ||
| / 10)), // Controls rate | ||
| Types.LONG), | ||
| WatermarkStrategy.noWatermarks(), | ||
| "ImpulseSource"); | ||
|
|
||
| for (String load : taskLoads) { | ||
| double maxLoad = Double.parseDouble(load); | ||
|
|
@@ -97,31 +143,6 @@ public static void main(String[] args) throws Exception { | |
| + ")"); | ||
| } | ||
|
|
||
| private static class ImpulseSource implements SourceFunction<Long> { | ||
| private final int maxSleepTimeMs; | ||
| volatile boolean canceled; | ||
|
|
||
| public ImpulseSource(int samplingInterval) { | ||
| this.maxSleepTimeMs = samplingInterval / 10; | ||
| } | ||
|
|
||
| @Override | ||
| public void run(SourceContext<Long> sourceContext) throws Exception { | ||
| while (!canceled) { | ||
| synchronized (sourceContext.getCheckpointLock()) { | ||
| sourceContext.collect(42L); | ||
| } | ||
| // Provide an impulse to keep the load simulation active | ||
| Thread.sleep(maxSleepTimeMs); | ||
| } | ||
| } | ||
|
|
||
| @Override | ||
| public void cancel() { | ||
| canceled = true; | ||
| } | ||
| } | ||
|
|
||
| private static class LoadSimulationFn extends RichFlatMapFunction<Long, Long> { | ||
|
|
||
| private final double maxLoad; | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
nit: either introduce this constant and use it in the calculations or use plain text in the comment.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Added the constant.