Apache NiFi并发任务机制及自定义处理器并行实现问询
Apache NiFi 并发并行处理机制及自定义处理器实现方式
核心并发逻辑
NiFi的并发处理由底层调度框架管控,当你将处理器的并发任务数设置为N(比如5)时,调度器会启动N个独立执行线程,每个线程都会重复调用处理器的onTrigger方法。每个线程调用onTrigger时,会独立从处理器输入队列中获取FlowFile并处理,以此实现并行处理。
自定义处理器的并发实现
对于继承自AbstractProcessor的自定义处理器,无需额外实现新方法——并发逻辑完全由NiFi调度框架接管,你只需确保两点:
onTrigger方法是线程安全的:不要在方法内使用非线程安全的共享变量,或对共享资源访问做同步处理- 每个
onTrigger调用仅处理单个(或少量)FlowFile:NiFi调度器会负责将不同FlowFile分配到不同onTrigger线程中,你只需在onTrigger里完成单FlowFile的业务逻辑即可
简言之,你编写的onTrigger方法就是并发执行的单元,NiFi会根据设置的并发数,启动多个线程同时调用该方法,每个线程处理各自的FlowFile实例。
内容的提问来源于stack exchange,提问作者Sathish
相关产品推荐
相关产品推荐

