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
Reversed Thread Scheduling
You mixed upsubscribeOn()andobserveOn():subscribeOn()defines which thread the Observable's creation/background work runs on (this should be a background thread likeSchedulers.io()for your frame processing)observeOn()defines which thread the Observer's callbacks (UI updates) run on (this should beAndroidSchedulers.mainThread())
Your current code runs the heavy frame processing on the main thread, which will block your UI and cause freezes/crashes.
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 onframeReady.acquire().Unclean Resource Handling
When breaking out of the loop (either becauseprocessingis 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, andprocessingare accessed safely across threads. Consider usingAtomicBooleanforprocessingif 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

