我正在测试 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)