Flink作业中多Source->Process->Sink处理逻辑的正式概念名称咨询
Answer to Your Question
Absolutely! The Source->Process->Sink pattern you're asking about has a widely recognized formal name in data processing: it's called a Data Pipeline (or specifically a Streaming Data Pipeline when working with real-time systems like Apache Flink).
In Flink's context, each independent flow that starts with a data source, goes through any number of processing operations (like transformations, aggregations, or filtering), and ends with a sink is exactly one such pipeline. Your code example shows two distinct data pipelines running within a single Flink job:
- Pipeline 1:
Source1→ (no explicit processing steps here) →Sink1 - Pipeline 2:
Source2→ (no explicit processing steps here) →Sink2
To elaborate a bit more:
- Even when you add intermediate processing logic (e.g.,
map,filter,keyBy, or window operations) between the source and sink, the entire end-to-end flow still qualifies as a data pipeline. - Flink lets you bundle multiple pipelines into a single job because they share the same
StreamExecutionEnvironment. This is handy when you want to manage related or unrelated streaming workloads in one deployment. - While Flink internally represents these flows as parts of a larger Operator DAG (Directed Acyclic Graph), "data pipeline" is the standard, high-level term developers use to describe this Source->Process->Sink pattern.
Here's your sample code for reference:
val env = StreamExecutionEnvironment.getExecutionEnvironment env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime) env.addSource(new Source1()).name("Source1").addSink(new Sink1()).name("Sink1").setParallelism(1) env.addSource(new Source2()).name("Source2").addSink(new Sink2()).name("Sink2").setParallelism(1)
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

