基于Flink-Kubernetes-Operator的作业多类异常排查求助
基于Flink-Kubernetes-Operator部署任务遇到的三类问题
环境信息
- Flink版本:1.14.4
- Flink-Kubernetes-Operator版本:release-1.4.0(commit:7fc23a1)
我通过Flink-Kubernetes-Operator部署Flink任务,配置了挂载PVC的Checkpoint目录,使用RocksDB状态后端并开启增量检查点,目前遇到以下三类问题:
问题1:重启服务时部分子任务阻塞在DEPLOYING/INITIALIZING状态
- 重启后部分子任务一直卡在DEPLOYING或INITIALIZING状态,无法继续推进,其余子任务正常处于RUNNING状态;删除任务后重新部署,偶尔能恢复正常。
问题2:增量检查点大小持续增长,疑似状态泄漏
- 已为自定义状态设置TTL,其余逻辑采用ReduceFunction,但增量检查点的大小仍在持续增长,不知道该从哪些方面排查状态泄漏的可能位置。
问题3:任务运行期间偶发阻塞,消费停滞且Checkpoint失败
- 任务运行过程中偶尔会出现类似阻塞的情况,消费完全停滞,同时Checkpoint持续失败;用arthas/jstack查看线程栈,发现存在线程阻塞现象,但每次阻塞的场景不同:
场景一线程栈:
[arthas@1]$ thread -b "Window(TumblingEventTimeWindows(60000), EventTimeTrigger, CommonWindowReduceFunction, PassThroughWindowFunction) -> Flat Map (22/72)#0" Id=157 BLOCKED on java.util.jar.JarFile@449cf4a0 owned by "Window(TumblingEventTimeWindows(60000), EventTimeTrigger, AppGroupTransactionEventApplyFunction) -> (Sink: Unnamed, Timestamps/Watermarks -> Flat Map) (31/32)#0" Id=151 at java.util.zip.ZipFile$ZipFileInputStream.read(ZipFile.java:719) - blocked on java.util.jar.JarFile@449cf4a0 at java.util.zip.ZipFile$ZipFileInflaterInputStream.fill(ZipFile.java:434) at java.util.zip.InflaterInputStream.read(InflaterInputStream.java:158) at sun.misc.Resource.getBytes(Resource.java:124) at java.net.URLClassLoader.defineClass(URLClassLoader.java:463) at java.net.URLClassLoader.access$100(URLClassLoader.java:74) at java.net.URLClassLoader$1.run(URLClassLoader.java:369) at java.net.URLClassLoader$1.run(URLClassLoader.java:363) at java.security.AccessController.doPrivileged(Native Method) at java.net.URLClassLoader.findClass(URLClassLoader.java:362) at org.apache.flink.util.ChildFirstClassLoader.loadClassWithoutExceptionHandling(ChildFirstClassLoader.java:71) at org.apache.flink.util.FlinkUserCodeClassLoader.loadClass(FlinkUserCodeClassLoader.java:48) - locked org.apache.flink.util.ChildFirstClassLoader@697c9014 <---- but blocks 93 other threads! at java.lang.ClassLoader.loadClass(ClassLoader.java:357) at org.apache.flink.runtime.execution.librarycache.FlinkUserCodeClassLoaders$SafetyNetWrapperClassLoader.loadClass(FlinkUserCodeClassLoaders.java:172) at java.lang.Class.forName0(Native Method) at java.lang.Class.forName(Class.java:348) at org.apache.flink.util.InstantiationUtil$ClassLoaderObjectInputStream.resolveClass(InstantiationUtil.java:78) at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1868) at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1751) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2042) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573) at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2287) at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2211) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2069) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573) at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2287) at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2211) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2069) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573) at java.io.ObjectInputStream.readObject(ObjectInputStream.java:431) at java.util.ArrayList.readObject(ArrayList.java:797) at sun.reflect.GeneratedMethodAccessor22.invoke(Unknown Source) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1170) at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2178) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2069) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573) at java.io.ObjectInputStream.readObject(ObjectInputStream.java:431) at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:617) at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:602) at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:589) at org.apache.flink.util.InstantiationUtil.readObjectFromConfig(InstantiationUtil.java:543) at org.apache.flink.streaming.api.graph.StreamConfig.getOutEdgesInOrder(StreamConfig.java:485) at org.apache.flink.streaming.runtime.tasks.StreamTask.createRecordWriters(StreamTask.java:1612) at org.apache.flink.streaming.runtime.tasks.StreamTask.createRecordWriterDelegate(StreamTask.java:1596) at org.apache.flink.streaming.runtime.tasks.StreamTask.<init>(StreamTask.java:376) at org.apache.flink.streaming.runtime.tasks.StreamTask.<init>(StreamTask.java:359) at org.apache.flink.streaming.runtime.tasks.StreamTask.<init>(StreamTask.java:332) at org.apache.flink.streaming.runtime.tasks.StreamTask.<init>(StreamTask.java:324) at org.apache.flink.streaming.runtime.tasks.StreamTask.<init>(StreamTask.java:314) at org.apache.flink.streaming.runtime.tasks.OneInputStreamTask.<init>(OneInputStreamTask.java:75) at sun.reflect.GeneratedConstructorAccessor37.newInstance(Unknown Source) at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45) at java.lang.reflect.Constructor.newInstance(Constructor.java:423) at org.apache.flink.runtime.taskmanager.Task.loadAndInstantiateInvokable(Task.java:1582) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:740) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575) at java.lang.Thread.run(Thread.java:748)
场景二线程栈:
"System Time Trigger for Window(TumblingEventTimeWindows(60000), EventTimeTrigger, CommonWindowReduceFunction, PassThroughWindowFunction) (32/48)#0" Id=232 TIMED_WAITING on java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject@783b47ad at sun.misc.Unsafe.park(Native Method) - waiting on java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject@783b47ad at java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:215) at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.awaitNanos(AbstractQueuedSynchronizer.java:2078) at java.util.concurrent.ScheduledThreadPoolExecutor$DelayedWorkQueue.take(ScheduledThreadPoolExecutor.java:1093) at java.util.concurrent.ScheduledThreadPoolExecutor$DelayedWorkQueue.take(ScheduledThreadPoolExecutor.java:809) at java.util.concurrent.ThreadPoolExecutor.getTask(ThreadPoolExecutor.java:1074) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1134) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ... "KeyedProcess -> Sink: CCLOUD_APP_SPAN_PORTRAIT_DETECT (24/32)#0" Id=134 BLOCKED on org.apache.flink.util.ChildFirstClassLoader@53624f6b owned by "Window(TumblingEventTimeWindows(60000), EventTimeTrigger, CommonWindowReduceFunction, PassThroughWindowFunction) (32/48)#0" Id=92 at java.lang.Class.getDeclaredFields0(Native Method) - blocked on org.apache.flink.util.ChildFirstClassLoader@53624f6b at java.lang.Class.privateGetDeclaredFields(Class.java:2583) at java.lang.Class.getDeclaredField(Class.java:2068) at java.io.ObjectStreamClass.getDeclaredSUID(ObjectStreamClass.java:1857) at java.io.ObjectStreamClass.access$700(ObjectStreamClass.java:79) at java.io.ObjectStreamClass$3.run(ObjectStreamClass.java:506) at java.io.ObjectStreamClass$3.run(ObjectStreamClass.java:494) at java.security.AccessController.doPrivileged(Native Method)
内容的提问来源于stack exchange,提问作者fidodosomething
相关产品推荐
相关产品推荐

