flink 1.14.3集群jar部署Recovery is suppressed by NoRestartBackoffTimeStrategy
flink程序在开发环境已经运行成功的情况下,部署到独立的flink集群(start-cluster)中,可能遇到不能正常运行的情况。
1. org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy
没有指定重启策略,在本地部署时,不需要指定重启策略。
可以通过下面的代码指定重启策略
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
                3, // 尝试重启的次数
                Time.of(10, TimeUnit.SECONDS) // 间隔
        ));
失败率重启
val env = ExecutionEnvironment.getExecutionEnvironment()
env.setRestartStrategy(RestartStrategies.failureRateRestart(
  3, // 一个时间段内的最大失败次数
  Time.of(5, TimeUnit.MINUTES), // 衡量失败次数的是时间段
  Time.of(10, TimeUnit.SECONDS) // 间隔
))
不重启
val env = ExecutionEnvironment.getExecutionEnvironment()
env.setRestartStrategy(RestartStrategies.noRestart())
2. getSerializableListState(Ljava/lang/String;)
在本地调试正常,上传到集群环境出现该错误提示:
org.apache.flink.api.common.state.OperatorStateStore.getSerializableListState(Ljava/lang/String;)
主要原因是版本不一致造成的,pom.xml文件中添加了scala 2.11的依赖包,集群环境是2.12版本,所以导致执行失败。
3. No ExecutorFactory found to execute the application
本地调试的时候,出现该错误,缺少依赖。
<flink.version>1.14.3</flink.version>
<dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients_2.12</artifactId>
            <version>${flink.version}</version>
        </dependency>
4. 可以打包在集群环境运行的参考pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>org.example</groupId>
    <artifactId>flink-shopstatis-demo</artifactId>
    <version>1.0</version>
    <properties>
        <maven.compiler.source>8</maven.compiler.source>
        <maven.compiler.target>8</maven.compiler.target>
        <flink.version>1.14.3</flink.version>
        <scala.version>2.12</scala.version>
        <project.build.scope>provided</project.build.scope>
    </properties>
    <dependencies>
        <!-- Flink 的 Java api -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>${flink.version}</version>
            <scope>${project.build.scope}</scope>
        </dependency>
        <!-- Flink Streaming 的 Java api -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java_${scala.version}</artifactId>
            <version>${flink.version}</version>
            <scope>${project.build.scope}</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-kafka_2.12</artifactId>
            <version>1.14.3</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-core</artifactId>
            <version>1.14.3</version>
            <scope>${project.build.scope}</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients_2.12</artifactId>
            <version>${flink.version}</version>
        </dependency>
        
    </dependencies>
    <build>
        <plugins>
            <!-- 编译插件 -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.6.0</version>
                <configuration>
                    <source>1.8</source>
                    <target>1.8</target>
                    <encoding>UTF-8</encoding>
                </configuration>
            </plugin>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>3.3.0</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies</descriptorRef>
                    </descriptorRefs>
                    <archive>
                        <manifest>
                            <!-- 可以设置jar包的入口类(可选) -->
                            <mainClass>statis.ShopDayGMVStatis</mainClass>
                        </manifest>
                    </archive>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>
打包时,对于集群环境中已经存在的jar,打包时采用provided,避免冲突。
 
                    版权声明:程序员胖胖胖虎阿 发表于 2022年9月8日 上午10:40。
转载请注明:flink 1.14.3集群jar部署Recovery is suppressed by NoRestartBackoffTimeStrategy | 胖虎的工具箱-编程导航
                    
            转载请注明:flink 1.14.3集群jar部署Recovery is suppressed by NoRestartBackoffTimeStrategy | 胖虎的工具箱-编程导航
相关文章
暂无评论...
