dalelane commented on code in PR #27579:
URL: https://github.com/apache/flink/pull/27579#discussion_r3645556469


##########
flink-runtime/src/main/java/org/apache/flink/runtime/execution/librarycache/BlobLibraryCacheManager.java:
##########
@@ -238,25 +238,41 @@ private UserCodeClassLoader getOrResolveClassLoader(
                 verifyIsNotReleased();
 
                 if (resolvedClassLoader == null) {
-                    boolean systemClassLoader =
-                            wrapsSystemClassLoader && libraries.isEmpty() && 
classPaths.isEmpty();
-                    resolvedClassLoader =
-                            new ResolvedClassLoader(
-                                    systemClassLoader
-                                            ? 
ClassLoader.getSystemClassLoader()
-                                            : createUserCodeClassLoader(
-                                                    jobId, applicationId, 
libraries, classPaths),
-                                    libraries,
-                                    classPaths,
-                                    systemClassLoader);
+                    resolvedClassLoader = createResolvedClassLoader(libraries, 
classPaths);
                 } else {
-                    resolvedClassLoader.verifyClassLoader(libraries, 
classPaths);
+                    try {
+                        resolvedClassLoader.verifyClassLoader(libraries, 
classPaths);
+                    } catch (IllegalStateException e) {
+                        LOG.warn(
+                                "Library cache entry for job {} has a 
classloader resolved with different "
+                                        + "library BLOBs than requested. This 
can happen during JobManager "
+                                        + "failover. Re-creating the 
classloader with the new blob keys.",
+                                jobId,
+                                e);
+                        resolvedClassLoader = 
createResolvedClassLoader(libraries, classPaths);

Review Comment:
   The previous resolvedClassLoader instance is overwritten without calling 
releaseClassLoader on it first. I can see that this is intentional - in the 
commit message, you explain that
   
   > The old classloader is intentionally not closed because in-flight tasks 
may still reference it.
   
   That does sound reasonable - it wouldn't be safe to call it here. But could 
we track it so it is closed later once no longer in use?
   
   If we keep the handle available from here, we can call releaseClassLoader() 
on that stale classloader as part of [the release method on line 
350](flink-runtime/src/main/java/org/apache/flink/runtime/execution/librarycache/BlobLibraryCacheManager.java:252).
   
   What do you think? 



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to