【问题标题】:Apache Flink, number of Task Slot vs env.setParallelismApache Flink,任务槽数与 env.setParallelism
【发布时间】: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


    【解决方案1】:

    通常每个插槽都会运行一个管道的并行实例。因此,作业的并行度与运行它所需的槽数相同。 (通过使用槽共享组,您可以将特定任务强制放入它们自己的槽中,这会增加所需槽的数量。)

    每个任务(包括一个或多个链接在一起的运算符)在一个 Java 线程中运行。

    任务管理器可以根据需要创建任意数量的插槽。典型配置每个插槽使用 1 个 CPU 内核,但对于处理要求高的管道,您可能希望每个插槽有 2 个或更多内核,而对于大部分空闲的管道,您可能会采用另一个方向并为每个内核配置多个插槽。

    在任务管理器中运行的所有任务/线程将简单地竞争任务管理器可以从托管它的机器或容器获取的 CPU 资源。

    所有状态对于使用它的一个操作员实例(任务)来说都是本地的,因此所有访问都发生在那个线程中。假设可能存在竞争条件的一个地方是 ProcessFunction 中的 onTimer 和 processElement 回调之间,但这些方法是同步的,因此您不必担心这一点。因为所有状态访问都是本地的,所以这会带来高吞吐量、低延迟和高可扩展性。

    在您的示例中,如果并行度为两个,那么您将有两个插槽独立地在数据的不同切片上执行相同的逻辑。如果他们正在使用状态,那么这将是由 Flink 管理的键分区状态,您可以将其视为分片键/值存储。

    对于时间窗口中的传感器数据,您完全不必担心多线程。 keyBy 将对数据进行分区,这样一个实例将处理某些传感器的所有事件和窗口,而另一个实例(假设有两个)将处理其余的。

    【讨论】:

    • 感谢您的回复。根据我的研究,我发现“每个taskManager都会有自己的数据(大数据将被划分到taskManager)”并且每个任务中的每个插槽都会并行执行数据?而且我还发现“一个TaskManager在同一个JVM进程中多线程执行它的任务”如果是真的,flink如何保证避免“竞争条件,共享访问数据”?如果理解正确的话,TaskManager 有 16 个插槽(机器有 16 个内核,比方说)意味着有一个数据(小型数据库等),其中 16 个插槽通过并行工作?
    • 有两点让这个简单:每个任务只有一个线程,每个状态只能从一个任务中访问。无法访问其他线程中的状态。 API 通过禁止共享数据访问来防止这些问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多