【发布时间】:2020-01-14 13:29:50
【问题描述】:
您能解释一下 Apache Flink v1.9 中任务槽和并行性之间的区别吗?
-
这是我目前的理解
- Flink 说TaskManager 是worker PROCESS。通常每台计算机应该有一个 TaskManager。
- 假设我有 3 台计算机,它们都有 16 个 CPU 内核。每台计算机都将是TaskManager。因此我将拥有 3 个 TaskManager
- 我曾想过如果一台电脑有16个cpu核心,那么TaskManager最多可以创建16个Task slot。因此那里有一个CPU隔离。但是 Flink 说 link => “请注意,这里没有 CPU 隔离;目前插槽仅分隔任务的托管内存。”
- 这意味着 16 个插槽 = 16 个线程?还有
numberOfSlot can be >= numberOfCpuCores?
如果任务槽意味着线程,这可能会导致“共享访问数据问题、竞争条件”等..?这是我的第一个问题。
- 第二个问题是我写到帖子开头的那个=> 任务槽和并行性之间的差异。我说的是 env.setparalellism(number)。
- 假设我的并行数 = 2
- 那么对于每个任务槽(线程或其他)将使用 2 个线程执行?
- 如果是,这可能会导致“共享访问数据问题、竞争条件”等?
- 如果不是,并行度是什么意思?
- 这是示例。在这个例子中,由于线程环境,我是否应该关心写
apply()method?:
public class AverageSensorReadings {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
int paralellism = env.getParallelism();
int maxParal = env.getMaxParallelism();
// ingest sensor stream
DataStream < SensorReading > sensorData = env
// SensorSource generates random temperature readings
.addSource(new SensorSource())
// assign timestamps and watermarks which are required for event time
.assignTimestampsAndWatermarks(new SensorTimeAssigner());
DataStream < SensorReading > avgTemp = sensorData
// convert Fahrenheit to Celsius using and inlined map function
.map(r -> new SensorReading(r.id, r.timestamp, (r.temperature - 32) * (5.0 / 9.0)))
// organize stream by sensor
.keyBy(r -> r.id)
// group readings in 1 second windows
.timeWindow(Time.seconds(4))
// compute average temperature using a user-defined function
.apply(new TemperatureAverager());
// print result stream to standard out
//avgTemp.print();
System.out.println("paral: " + paralellism + " max paral: " + maxParal);
// execute application
env.execute("Compute average sensor temperature");
}
public static class TemperatureAverager extends RichWindowFunction < SensorReading, SensorReading, String, TimeWindow > {
/**
* apply() is invoked once for each window.
*
* @param sensorId the key (sensorId) of the window
* @param window meta data for the window
* @param input an iterable over the collected sensor readings that were assigned to the window
* @param out a collector to emit results from the function
*/
@Override
public void apply(String sensorId, TimeWindow window, Iterable < SensorReading > input, Collector < SensorReading > out) {
System.out.println("APPLY FUNCTION START POINT");
System.out.println("sensorId: " + sensorId + "\n");
// compute the average temperature
int cnt = 0;
double sum = 0.0;
for (SensorReading r: input) {
System.out.println("collected item: " + r);
cnt++;
sum += r.temperature;
}
double avgTemp = sum / cnt;
System.out.println("APPLY FUNCTION END POINT");
System.out.println("----------------------------\n\n");
// emit a SensorReading with the average temperature
out.collect(new SensorReading(sensorId, window.getEnd(), avgTemp));
}
}
}
【问题讨论】:
标签: java stream apache-flink flink-streaming