[FLINK-3313] [kafka] Fix TypeInformationSerializationSchema usage in LegacyFetcher
The LegacyFetcher used the given KeyedDeserializationSchema across multiple threads even though it is not thread-safe. This commit fixes the problem by cloning the KeyedDeserializationSchema before giving it to the SimpleConsumerThread. Add clone method for Java serializable objects to InstantiationUtil This closes #1577.
Showing
想要评论请 注册 或 登录