Hazelcast Jet管道Sink中注入Spring依赖的可行方案探讨
在Hazelcast Jet Pipeline的Update函数中注入Spring依赖的可行方案
我们的场景是在Hazelcast Jet Pipeline的Sinks.mapWithUpdating的updateFn中注入Spring Bean,现有基于@SpringAware的自定义SpringInitializationAware实现因未等待上下文刷新、且Jet内部未触发ManagedContext::initialize导致依赖注入失败,以下是几种可行方案:
方案1:使用Spring上下文静态持有者获取Bean
通过静态类持有Spring应用上下文,在updateFn中直接从上下文获取依赖,方式简单直接:
// 静态上下文持有者类 @Component public class SpringContextHolder implements ApplicationContextAware { private static ApplicationContext context; @Override public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { context = applicationContext; } public static <T> T getBean(Class<T> beanClass) { return context.getBean(beanClass); } } // 修改后的Updater实现 public static class Updater implements BiFunctionEx<A, B, C> { @Override public C applyEx(A a, B b) { OurDependency dependency = SpringContextHolder.getBean(OurDependency.class); // 使用dependency处理业务逻辑 } }
需确保SpringContextHolder被Spring扫描初始化,且上下文在Jet任务启动前已完全刷新。
方案2:自定义Jet的ManagedContext
利用Hazelcast的ManagedContext扩展机制,实现Spring-aware的上下文初始化逻辑,让Jet在反序列化updateFn时自动完成依赖注入:
- 实现Spring-aware的ManagedContext:
public class SpringManagedContext implements ManagedContext { private final ApplicationContext applicationContext; public SpringManagedContext(ApplicationContext applicationContext) { this.applicationContext = applicationContext; } @Override public Object initialize(Object obj) { applicationContext.getAutowireCapableBeanFactory().autowireBean(obj); return obj; } }
- 配置Hazelcast Jet实例使用该ManagedContext:
Config config = new Config(); config.setManagedContext(new SpringManagedContext(springApplicationContext)); HazelcastInstance jetInstance = HazelcastJet.newJetInstance(config);
- 给
Updater类添加Spring注解(无需继承自定义的SpringInitializationAware):
@Component public static class Updater implements BiFunctionEx<A, B, C> { @Autowired private OurDependency dependency; @Override public C applyEx(A a, B b) { // 使用dependency处理业务逻辑 } }
方案3:将依赖作为参数传递到Pipeline
如果依赖是无状态或可序列化的,可在构建Pipeline时直接将依赖实例传递给updateFn:
// 在Spring上下文初始化后获取依赖实例 OurDependency dependency = springApplicationContext.getBean(OurDependency.class); // 构建Pipeline时传递依赖 Pipeline p = Pipeline.create(); p.readFrom(TestSources.items("the", "quick", "brown", "fox")) .writeTo(Sinks.mapWithUpdating( imap, toKeyFn, (a, b) -> { // 使用传入的dependency处理逻辑 return dependency.process(a, b); } ));
注意:若OurDependency是有状态的,需确保它可序列化,或通过@Scope("singleton")保证实例一致性。
内容的提问来源于stack exchange,提问作者Eugene
相关产品推荐
相关产品推荐

