小编Luc*_*lix的帖子

集成 - Apache Flink + Spring Boot

我正在测试 Apache Flink 和 Spring Boot 之间的集成,在 IDE 上运行它们很好,但是当我尝试在 Apache Flink 集群上运行时,我有一个与 ClassLoader 相关的异常。

这些课程非常简单:

BootFlink 应用程序

@SpringBootApplication
@ComponentScan("com.example.demo")
public class BootFlinkApplication {

    public static void main(String[] args) {
        System.out.println("some test");
        SpringApplication.run(BootFlinkApplication.class, args);
    }
}
Run Code Online (Sandbox Code Playgroud)

FlinkTest

@Service
public class FlinkTest {
    @PostConstruct
    public void init() {
        StreamExecutionEnvironment see = StreamExecutionEnvironment.getExecutionEnvironment();

        see.fromElements(1, 2, 3, 4)
            .filter(new RemoveNumber3Filter()).print();

        try {
            see.execute();
        } catch (Exception e) {
            System.out.println("Error executing flink job: " + e.getMessage());
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

RemoveNumber3Filter

public class RemoveNumber3Filter implements FilterFunction<Integer> …
Run Code Online (Sandbox Code Playgroud)

spring-boot apache-flink

6
推荐指数
1
解决办法
9151
查看次数

标签 统计

apache-flink ×1

spring-boot ×1