【问题标题】:Why does Apache Flink stall in method org.apache.flink.api.java.typeutils.runtime.kryo.Serializers.getContainedGenericTypes?为什么 Apache Flink 在方法 org.apache.flink.api.java.typeutils.runtime.kryo.Serializers.getContainedGenericTypes 中停止?
【发布时间】:2017-03-17 15:19:13
【问题描述】:

在实现我的算法时,我在 Apache Flink 中使用 for 循环创建了一长串运算符。从方法中的一些长度处理停止开始 org.apache.flink.api.java.typeutils.runtime.kryo.Serializers.getContainedGenericTypes 很久才实际处理。如何解释这种现象?如何解决它以减少此方法时间?

【问题讨论】:

  • 您在流中使用哪些数据类型?您的 Kryo 类型可能未注册。
  • 我正在使用带有原语的案例类类型,例如案例类 Cell (i:Int.j:Int,v1:Int,v2:Int)。我正在开发用于 DataSet[Cell] 批处理的系统。
  • @rmetzger 我应该以某种方式显式注册这些类型吗?
  • 不,你不需要注册类型。这就是 Serializers 类正在做的事情。

标签: scala apache-flink kryo


【解决方案1】:

Serializers.getContainedGenericTypes() 方法仅在您的DataSet 应用程序的计划创建期间调用。

设置ExecutionConfig.disableAutoTypeRegistration() 将禁用此注册。

我假设您在本地运行您的 Flink 应用程序而没有大量数据。通常,计划创建只占用可用 CPU 时间的一小部分,而实际处理会占用大部分时间。

【讨论】:

  • 是的,假设是正确的。 Flink 在这次运行中应该处理几兆字节。我希望在我的研究中处理几个 TB 的集群。禁用自动类型注册后,我遇到了类似的麻烦。我在分析器中得到了org.apache.flink.optimizer.traversals ~50%、org.apache.flink.optimizer.plan ~32%、java.util ~17% 的 CPU 时间,并且几个小时没有响应。
  • 处理几兆字节,你不需要一个大的分布式处理框架。 Java 集合将完成这项工作。一旦要处理几 TB,优化器或类型注册将不再占主导地位。
  • 问题是它挂了。图表未启动。这不是一个简单的开销。我想运行一个生成 6 层的算法。在此我试图展示图表gist.github.com/protsenkovi/3f0cba82978c8ea41bd191cd9a1b1714 的复杂性。
猜你喜欢
  • 2020-05-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-10-20
  • 1970-01-01
相关资源
最近更新 更多