AHeise commented on code in PR #29356:
URL: https://github.com/apache/flink/pull/29356#discussion_r4165517921


##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/generated/CompileUtils.java:
##########
@@ -95,8 +107,32 @@ public static <T> Class<T> compile(ClassLoader cl, String 
name, String code) {
         }
     }
 
+    @SuppressWarnings("unchecked")
     private static <T> Class<T> doCompile(ClassLoader cl, String name, String 
code) {
         checkNotNull(cl, "Classloader must not be null.");
+        // Generated code only references types that resolve the same way 
under every classloader,
+        // so the bytecode is cooked once per code and defined into cl.
+        final Map<String, byte[]> byteCodes = sharedByteCode(cl, name, code);
+        try {
+            return (Class<T>) new ByteArrayClassLoader(cl, 
byteCodes).loadClass(name);
+        } catch (ClassNotFoundException e) {
+            throw new FlinkRuntimeException("Can not load class " + name, e);
+        }
+    }
+
+    private static Map<String, byte[]> sharedByteCode(ClassLoader cl, String 
name, String code) {
+        try {
+            return BYTECODE_CACHE.get(code, () -> cook(cl, name, code));

Review Comment:
   It's not a general Java concern, it's exactly what our codegen emits. 
`StructuredObjectConverter#castExpr` wraps both `int` and `Integer` in 
`(java.lang.Integer)`, so the string is identical while Janino binds 
`getAge()I` vs `getAge()Ljava/lang/Integer;`. Setters have the same issue via 
`setterExpr`. The converter name doesn't help: `nextUniqueClassId` is per 
planner JVM, so two `flink run` processes both start at 
`com$x$Pojo$0$Converter`. Stateless UDFs are no better, since 
`functionIdentifier()` is the bare class name.
   
   Standalone repro with Janino 3.1.12: compile the converter shape against the 
v1 `Pojo` and define it into the v2 loader, as `BYTECODE_CACHE` does:
   ```
   fresh cook in B -> 7
   shared A-bytes in B -> java.lang.NoSuchMethodError: 'int com.x.Pojo.getAge()'
   ```
   
   <details><summary>Repro.java</summary>
   
   ```java
   import java.util.Map;
   import org.codehaus.janino.SimpleCompiler;
   import org.codehaus.commons.compiler.util.reflect.ByteArrayClassLoader;
   
   public class Repro {
       static Map<String, byte[]> compile(String code, ClassLoader parent) 
throws Exception {
           SimpleCompiler c = new SimpleCompiler();
           c.setParentClassLoader(parent);
           c.cook(code);
           return c.getBytecodes();
       }
   
       public static void main(String[] args) throws Exception {
           ClassLoader app = Repro.class.getClassLoader();
           // job A jar: int getAge(); job B jar: Integer getAge()
           ClassLoader jobA = new ByteArrayClassLoader(compile(
               "package com.x; public class Pojo { public int getAge() { return 
42; } }", app), app);
           ClassLoader jobB = new ByteArrayClassLoader(compile(
               "package com.x; public class Pojo { public Integer getAge() { 
return 7; } }", app), app);
   
           // shape emitted by StructuredObjectConverter.getterExpr/castExpr 
for both int and Integer
           String gen = "public class Conv { public Object get(Object o) {"
               + " final com.x.Pojo external = (com.x.Pojo) o;"
               + " return ((java.lang.Integer) external.getAge()); } }";
   
           Map<String, byte[]> bytesA = compile(gen, jobA);  // cooked once, as 
with BYTECODE_CACHE
           Object pojoB = 
jobB.loadClass("com.x.Pojo").getDeclaredConstructor().newInstance();
   
           Class<?> fresh = new ByteArrayClassLoader(compile(gen, jobB), 
jobB).loadClass("Conv");
           System.out.println("fresh cook in B -> " + fresh.getMethod("get", 
Object.class)
               .invoke(fresh.getDeclaredConstructor().newInstance(), pojoB));
   
           Class<?> shared = new ByteArrayClassLoader(bytesA, 
jobB).loadClass("Conv");
           try {
               System.out.println("shared A-bytes in B -> " + 
shared.getMethod("get", Object.class)
                   .invoke(shared.getDeclaredConstructor().newInstance(), 
pojoB));
           } catch (java.lang.reflect.InvocationTargetException e) {
               System.out.println("shared A-bytes in B -> " + e.getCause());
           }
       }
   }
   ```
   </details>
   
   Realistic trigger: upgrading a job on a session cluster within the 5-minute 
window. It also moves the failure: before, a mismatch surfaced at instantiation 
as "Table program cannot be compiled"; now `defineClass` succeeds and the 
`NoSuchMethodError`/`NoClassDefFoundError` hits on first record.
   
   Same text only implies same bytecode if all referenced types resolve to the 
same classes. I'd guard the reuse on that: record what Janino resolved during 
the cook and reuse only if the new loader returns the same `Class` instances. 
That keeps all hits for Flink/JDK-only code.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/generated/CompileUtils.java:
##########
@@ -95,8 +107,32 @@ public static <T> Class<T> compile(ClassLoader cl, String 
name, String code) {
         }
     }
 
+    @SuppressWarnings("unchecked")
     private static <T> Class<T> doCompile(ClassLoader cl, String name, String 
code) {
         checkNotNull(cl, "Classloader must not be null.");
+        // Generated code only references types that resolve the same way 
under every classloader,
+        // so the bytecode is cooked once per code and defined into cl.
+        final Map<String, byte[]> byteCodes = sharedByteCode(cl, name, code);
+        try {
+            return (Class<T>) new ByteArrayClassLoader(cl, 
byteCodes).loadClass(name);
+        } catch (ClassNotFoundException e) {
+            throw new FlinkRuntimeException("Can not load class " + name, e);
+        }
+    }
+
+    private static Map<String, byte[]> sharedByteCode(ClassLoader cl, String 
name, String code) {
+        try {
+            return BYTECODE_CACHE.get(code, () -> cook(cl, name, code));
+        } catch (ExecutionException | UncheckedExecutionException e) {
+            Throwable cause = e.getCause();
+            if (cause instanceof RuntimeException) {
+                throw (RuntimeException) cause;
+            }
+            throw new FlinkRuntimeException(cause);

Review Comment:
   Minor: `InvalidProgramException` gets unwrapped here and then wrapped again 
in `FlinkRuntimeException` by `compile()`, so the cause chain differs from 
before. Harmless for `GeneratedClass` (catches `Throwable`), but double-check 
no test asserts on the message chain.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/generated/CompileUtils.java:
##########
@@ -95,8 +107,32 @@ public static <T> Class<T> compile(ClassLoader cl, String 
name, String code) {
         }
     }
 
+    @SuppressWarnings("unchecked")
     private static <T> Class<T> doCompile(ClassLoader cl, String name, String 
code) {
         checkNotNull(cl, "Classloader must not be null.");
+        // Generated code only references types that resolve the same way 
under every classloader,

Review Comment:
   This only holds for types loaded parent-first (JDK, `org.apache.flink.*`). 
UDFs, POJOs, RAW classes and user serializers resolve child-first from each 
job's jar, and Janino bakes overload choice, constant inlining and verifier 
assignability against the first loader. See the thread on L125 for a concrete 
`NoSuchMethodError` from `StructuredObjectConverter`.
   
   Proposal: during `cook`, wrap `cl` in a loader that records `name → Class` 
for every type Janino resolves; on reuse, take the cached bytes only if 
`cl.loadClass(n)` returns the same `Class` for all of them, otherwise cook 
again. Bytecode is then provably identical, and all hits for Flink/JDK-only 
code (cast executors, reducers, equalisers, projections) remain. Keep the 
recorded classes weak so the first loader isn't pinned.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/generated/CompileUtils.java:
##########
@@ -171,6 +202,25 @@ public static ExpressionEvaluator compileExpression(
         }
     }
 
+    /** Defines classes from cached bytecode, delegating referenced types to 
the parent. */
+    private static final class ByteArrayClassLoader extends ClassLoader {
+        private final Map<String, byte[]> byteCodes;
+
+        ByteArrayClassLoader(ClassLoader parent, Map<String, byte[]> 
byteCodes) {
+            super(parent);
+            this.byteCodes = byteCodes;
+        }
+
+        @Override
+        protected Class<?> findClass(String className) throws 
ClassNotFoundException {
+            final byte[] byteCode = byteCodes.get(className);
+            if (byteCode == null) {
+                throw new ClassNotFoundException(className);
+            }
+            return defineClass(className, byteCode, 0, byteCode.length);
+        }
+    }
+
     /** Class to use as key for the {@link #COMPILED_CLASS_CACHE}. */
     private static class ClassKey {
         private final int classLoaderId;

Review Comment:
   Not introduced here, but `cl.hashCode()` is the identity hash (no Flink 
classloader overrides it), and it's not unique. Two colliding loaders compiling 
the same code get a `Class` defined under the other job's loader, which links 
against the wrong jar or a closed loader. With deterministic names, "same code" 
is the common case. Now that the cook is shared, the class cache can afford an 
exact identity check, e.g. a `WeakReference<ClassLoader>` in `ClassKey` 
compared with `ref.get() == other`. Fine as a follow-up ticket.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/generated/CompileUtils.java:
##########
@@ -171,6 +202,25 @@ public static ExpressionEvaluator compileExpression(
         }
     }
 
+    /** Defines classes from cached bytecode, delegating referenced types to 
the parent. */
+    private static final class ByteArrayClassLoader extends ClassLoader {

Review Comment:
   Janino already ships this: 
`org.codehaus.commons.compiler.util.reflect.ByteArrayClassLoader(Map<String, 
byte[]>, ClassLoader)`. Same thing plus protection-domain handling. Any reason 
not to reuse it?



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/generated/CompileUtils.java:
##########
@@ -67,10 +70,19 @@ public final class CompileUtils {
                     .softValues()
                     .build();
 
+    /** Bytecode by code, cooked once and shared across classloaders. */
+    static final Cache<String, Map<String, byte[]>> BYTECODE_CACHE =

Review Comment:
   Do you have numbers? `defineClass` plus verification still happens per 
loader; only the cook is saved, and Metaspace is unchanged. The 
`COMPILED_CLASS_CACHE` Javadoc cites Metaspace as its motivation, so it would 
help to state what this cache targets. I'd guess planner hosts with a loader 
per session (`ExpressionReducer`, cast rules) and session clusters with short 
jobs. A small JMH comparing cook vs define for a typical split Calc/Agg would 
settle it.



##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/generated/CompileUtilsTest.java:
##########
@@ -39,6 +40,39 @@ void before() {
         // cleanup cached class before tests
         CompileUtils.COMPILED_CLASS_CACHE.invalidateAll();
         CompileUtils.COMPILED_EXPRESSION_CACHE.invalidateAll();
+        CompileUtils.BYTECODE_CACHE.invalidateAll();
+    }
+
+    @Test
+    void testSameSourceIsCompiledOnceAcrossClassLoaders() throws Exception {

Review Comment:
   Please add a negative case: same code, second loader whose probe has a 
different signature (e.g. the int/Integer getter from the L125 thread), 
asserting a recompile. Also worth covering code with an inner/anonymous class 
to exercise the multi-entry bytecode map.



-- 
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