You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何控制RxJava无限发射Observable的停止/恢复/终止?——Android RFID扫描场景技术问询

解决方案:基于RxJava实现RFID扫描的暂停/恢复/终止

针对你的需求——无需从头启动的扫描暂停/恢复、Fragment切换时的资源管理,我们可以通过原子状态控制+RxJava订阅生命周期管理来解决,同时避开takeUntil的终止性问题。下面是具体实现思路和代码:

一、先理清核心问题

你当前的代码有两个关键痛点:

  1. 每次点击按钮都会创建新的Observable和订阅,旧的扫描循环会在后台继续运行,导致内存泄漏和重复扫描;
  2. takeUntil会直接终止Observable序列,一旦触发就无法恢复扫描,不符合你“无需从头开始”的需求。

所以我们需要让扫描循环本身支持暂停/继续,同时统一管理订阅的生命周期。

二、重构Scanner单例类

我们在单例中加入原子状态变量来控制扫描的启动、暂停、停止,同时复用扫描的Observable和订阅:

class Scanner {
    private static Scanner instance;
    // 控制扫描是否正在运行
    private final AtomicBoolean isScanning = new AtomicBoolean(false);
    // 控制扫描是否处于暂停状态
    private final AtomicBoolean isPaused = new AtomicBoolean(false);
    // 管理扫描订阅的Disposable
    private Disposable scanDisposable;
    private ObservableEmitter<Product> currentEmitter;

    private Scanner() {}

    public static Scanner getInstance() {
        if (instance == null) {
            synchronized (Scanner.class) {
                if (instance == null) {
                    instance = new Scanner();
                }
            }
        }
        return instance;
    }

    // 启动或恢复扫描,传入UI层的Observer
    public void startOrResumeScan(Observer<Product> observer) {
        if (isScanning.get()) {
            // 正在扫描但处于暂停状态,直接恢复
            if (isPaused.get()) {
                isPaused.set(false);
            }
            return;
        }

        // 首次启动扫描
        isScanning.set(true);
        isPaused.set(false);

        scanDisposable = Observable.create((ObservableEmitter<Product> emitter) -> {
            currentEmitter = emitter;
            while (isScanning.get()) {
                // 暂停时进入循环等待,避免CPU空转
                while (isPaused.get()) {
                    Thread.sleep(100);
                    // 检查是否已停止扫描,防止无限等待
                    if (!isScanning.get()) break;
                }
                if (!isScanning.get()) break;

                // 执行RFID读取
                UHFTAGInfo TAG = RFIDReader.inventorySingleTag();
                if (TAG != null && !emitter.isDisposed()) {
                    emitter.onNext(new Product(TAG.getEPC()));
                }
                // 添加短延迟,降低CPU占用
                Thread.sleep(50);
            }
            emitter.onComplete();
        })
        .subscribeOn(Schedulers.io())
        .observeOn(AndroidSchedulers.mainThread())
        .subscribe(observer);
    }

    // 暂停扫描
    public void pauseScan() {
        if (isScanning.get() && !isPaused.get()) {
            isPaused.set(true);
        }
    }

    // 彻底停止扫描,释放资源
    public void stopScan() {
        isScanning.set(false);
        isPaused.set(false);
        if (scanDisposable != null && !scanDisposable.isDisposed()) {
            scanDisposable.dispose();
        }
        currentEmitter = null;
    }

    // 对外暴露状态,供UI层判断
    public boolean isScanning() {
        return isScanning.get();
    }

    public boolean isPaused() {
        return isPaused.get();
    }
}

三、UI层ScanFragment的适配

在Fragment中用CompositeDisposable统一管理订阅,处理按钮点击和生命周期事件:

public class ScanFragment extends Fragment {
    private Button scanControlButton;
    private List<Product> listProducts = new ArrayList<>();
    private ProductAdapter adapter;
    private final CompositeDisposable compositeDisposable = new CompositeDisposable();
    // 复用同一个Observer,避免重复创建
    private final Observer<Product> scanObserver = new Observer<Product>() {
        @Override
        public void onSubscribe(@NonNull Disposable d) {
            compositeDisposable.add(d);
        }

        @Override
        public void onNext(@NonNull Product product) {
            listProducts.add(product);
            adapter.notifyDataSetChanged();
        }

        @Override
        public void onError(@NonNull Throwable e) {
            Toast.makeText(getContext(), "扫描出错: " + e.getMessage(), Toast.LENGTH_SHORT).show();
            Scanner.getInstance().stopScan();
            updateButtonText();
        }

        @Override
        public void onComplete() {
            updateButtonText();
        }
    };

    @Override
    public View onCreateView(LayoutInflater inflater, ViewGroup container, Bundle savedInstanceState) {
        View view = inflater.inflate(R.layout.fragment_scan, container, false);
        scanControlButton = view.findViewById(R.id.scan_control_button);
        // 初始化RecyclerView和adapter...

        scanControlButton.setOnClickListener(v -> {
            Scanner scanner = Scanner.getInstance();
            if (scanner.isScanning()) {
                if (scanner.isPaused()) {
                    scanner.startOrResumeScan(scanObserver);
                } else {
                    scanner.pauseScan();
                }
            } else {
                scanner.startOrResumeScan(scanObserver);
            }
            updateButtonText();
        });

        return view;
    }

    // 更新按钮文本,反馈当前扫描状态
    private void updateButtonText() {
        Scanner scanner = Scanner.getInstance();
        if (scanner.isScanning()) {
            scanControlButton.setText(scanner.isPaused() ? "恢复扫描" : "暂停扫描");
        } else {
            scanControlButton.setText("开始扫描");
        }
    }

    @Override
    public void onPause() {
        super.onPause();
        // 切换到其他Fragment时暂停扫描,切回来可以直接恢复
        Scanner.getInstance().pauseScan();
        updateButtonText();
    }

    @Override
    public void onDestroy() {
        super.onDestroy();
        // 销毁Fragment时彻底停止扫描,释放所有资源
        Scanner.getInstance().stopScan();
        compositeDisposable.dispose();
    }
}

四、关键细节说明

  1. 原子变量的作用:AtomicBoolean保证多线程下的状态安全(扫描在IO线程,UI操作在主线程),避免竞态条件;
  2. 暂停机制:暂停时让扫描循环进入短睡眠等待,而不是终止Observable,这样恢复时可以直接继续扫描,无需重新初始化;
  3. 订阅管理:CompositeDisposable统一管理所有订阅,在Fragment销毁时dispose,彻底避免内存泄漏;
  4. 关于takeUntil的疑问:takeUntil确实会终止整个Observable序列,一旦触发就无法恢复,因此它更适合实现一次性的终止操作(比如用户明确点击“停止”且不再恢复的场景),不适合你的暂停/恢复需求。

内容的提问来源于stack exchange,提问作者tur11ng

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 00:24:05