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

Android中AsyncTask转RxJava遇问题,求代码修正方案

Hey there! Let's figure out why your RxJava code isn't working properly—there are a few key issues in your implementation that we can fix step by step.

Key Issues in Your Current RxJava Code

  1. Reversed Thread Scheduling
    You mixed up subscribeOn() and observeOn():

    • subscribeOn() defines which thread the Observable's creation/background work runs on (this should be a background thread like Schedulers.io() for your frame processing)
    • observeOn() defines which thread the Observer's callbacks (UI updates) run on (this should be AndroidSchedulers.mainThread())
      Your current code runs the heavy frame processing on the main thread, which will block your UI and cause freezes/crashes.
  2. Missing Disposal Safety Checks
    You don't check if the emitter has been disposed before continuing your loop. If the subscription is cancelled (e.g., when the Activity/Fragment is destroyed), your loop might keep running or get stuck blocking on frameReady.acquire().

  3. Unclean Resource Handling
    When breaking out of the loop (either because processing is false or the emitter is disposed), you should ensure you're not leaving resources locked or hanging.

Fixed RxJava Implementation

Here's the corrected code with explanations in comments:

// Store this Disposable as a member variable so you can dispose it later (e.g., in onDestroy)
private Disposable processingDisposable;

public void workEngine() {
    // Cancel any existing processing before starting a new one
    if (processingDisposable != null && !processingDisposable.isDisposed()) {
        processingDisposable.dispose();
    }

    processingDisposable = Observable.create((ObservableOnSubscribe<RecognitionResult>) emitter -> {
        try {
            while (!emitter.isDisposed() && processing) { // Check disposal + processing flag
                frameReady.acquire(); // Wait for frame

                // Double-check after acquiring the lock (processing might have been set to false while waiting)
                if (!processing || emitter.isDisposed()) {
                    frameReady.release(); // Don't leave the lock acquired!
                    break;
                }

                Camera.Size size = camera.getParameters().getPreviewSize();
                RecognitionResult result;

                switch (mAngle) {
                    case 0:
                        result = session.ProcessYUVSnapshot(mData, size.width, size.height, ImageOrientation.Landscape);
                        break;
                    case 180:
                        result = session.ProcessYUVSnapshot(mData, size.width, size.height, ImageOrientation.InvertedLandscape);
                        break;
                    case 270:
                        result = session.ProcessYUVSnapshot(mData, size.width, size.height, ImageOrientation.InvertedPortrait);
                        break;
                    default:
                        result = session.ProcessYUVSnapshot(mData, size.width, size.height, ImageOrientation.Portrait);
                }

                // Only send the result if the emitter is still active
                if (!emitter.isDisposed()) {
                    emitter.onNext(result);
                }
            }
        } catch (InterruptedException e) {
            // Handle interruptions (e.g., when disposing the subscription)
            Thread.currentThread().interrupt();
        } catch (Exception e) {
            String message = "Error while processing frame: " + e.toString();
            Log.d("smartid", message);
            callback.error(message);
            
            // Send error to observer if emitter is active
            if (!emitter.isDisposed()) {
                emitter.onError(e);
            }
        } finally {
            // Ensure we release any remaining locks if we exit the loop
            if (frameReady.hasQueuedThreads()) {
                frameReady.release();
            }
            // Complete the observable when done
            if (!emitter.isDisposed()) {
                emitter.onComplete();
            }
        }
    })
    .subscribeOn(Schedulers.io()) // Run background processing on IO thread
    .observeOn(AndroidSchedulers.mainThread()) // Run UI updates on main thread
    .subscribe(new Observer<RecognitionResult>() {
        @Override
        public void onSubscribe(Disposable d) {
            // Optional: You can store the disposable here too, but we already have it as a member
        }

        @Override
        public void onNext(RecognitionResult recognitionResult) {
            // Update UI safely on main thread
            draw.showResult(recognitionResult);
            draw.invalidate();
            callback.recognized(recognitionResult);
            frameWaiting.release();
        }

        @Override
        public void onComplete() {
            Log.d("smartid", "Frame processing completed");
        }

        @Override
        public void onError(Throwable e) {
            e.printStackTrace();
            callback.error("Processing failed: " + e.getMessage());
        }
    });
}

// Don't forget to dispose the subscription when your component is destroyed!
// e.g., in Activity.onDestroy() or Fragment.onDestroyView()
@Override
protected void onDestroy() {
    super.onDestroy();
    if (processingDisposable != null && !processingDisposable.isDisposed()) {
        processingDisposable.dispose();
        processing = false; // Ensure the loop exits
    }
}

Additional Notes for Android Beginners

  • Disposable Management: Always dispose subscriptions when your Activity/Fragment is destroyed to prevent memory leaks and background work from continuing unnecessarily.
  • Thread Safety: Make sure variables like mData, camera, and processing are accessed safely across threads. Consider using AtomicBoolean for processing if multiple threads are modifying it.
  • Exception Handling: In RxJava, unhandled exceptions can terminate the observable, so we added proper error handling and ensure we don't send events to a disposed emitter.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:11:20