This is an automated email from the ASF dual-hosted git repository.

forwardxu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 60b62fc25c [MINOR] retain avro's namespace (#6783)
60b62fc25c is described below

commit 60b62fc25c0b7d34ff011a16c68b5d699c8e6bbc
Author: Yann Byron <[email protected]>
AuthorDate: Sun Sep 25 05:19:30 2022 +0800

    [MINOR] retain avro's namespace (#6783)
---
 .../convert/AvroInternalSchemaConverter.java       | 31 +++++++++++++---------
 .../schema/utils/AvroSchemaEvolutionUtils.java     |  4 ++-
 2 files changed, 22 insertions(+), 13 deletions(-)

diff --git 
a/hudi-common/src/main/java/org/apache/hudi/internal/schema/convert/AvroInternalSchemaConverter.java
 
b/hudi-common/src/main/java/org/apache/hudi/internal/schema/convert/AvroInternalSchemaConverter.java
index 360134f928..d941b27328 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/internal/schema/convert/AvroInternalSchemaConverter.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/internal/schema/convert/AvroInternalSchemaConverter.java
@@ -50,8 +50,12 @@ public class AvroInternalSchemaConverter {
    * @param tableName the record name.
    * @return an avro Schema.
    */
+  public static Schema convert(InternalSchema internalSchema, String 
tableName, String namespace) {
+    return buildAvroSchemaFromInternalSchema(internalSchema, tableName, 
namespace);
+  }
+
   public static Schema convert(InternalSchema internalSchema, String 
tableName) {
-    return buildAvroSchemaFromInternalSchema(internalSchema, tableName);
+    return buildAvroSchemaFromInternalSchema(internalSchema, tableName, "");
   }
 
   /**
@@ -241,7 +245,7 @@ public class AvroInternalSchemaConverter {
    */
   public static Schema buildAvroSchemaFromType(Type type, String recordName) {
     Map<Type, Schema> cache = new HashMap<>();
-    return visitInternalSchemaToBuildAvroSchema(type, cache, recordName);
+    return visitInternalSchemaToBuildAvroSchema(type, cache, recordName, "");
   }
 
   /**
@@ -251,9 +255,9 @@ public class AvroInternalSchemaConverter {
    * @param recordName the record name
    * @return a Avro schema match hudi internal schema.
    */
-  public static Schema buildAvroSchemaFromInternalSchema(InternalSchema 
schema, String recordName) {
+  public static Schema buildAvroSchemaFromInternalSchema(InternalSchema 
schema, String recordName, String namespace) {
     Map<Type, Schema> cache = new HashMap<>();
-    return visitInternalSchemaToBuildAvroSchema(schema.getRecord(), cache, 
recordName);
+    return visitInternalSchemaToBuildAvroSchema(schema.getRecord(), cache, 
recordName, namespace);
   }
 
   /**
@@ -264,13 +268,15 @@ public class AvroInternalSchemaConverter {
    * @param recordName the record name
    * @return a Avro schema match this type
    */
-  private static Schema visitInternalSchemaToBuildAvroSchema(Type type, 
Map<Type, Schema> cache, String recordName) {
+  private static Schema visitInternalSchemaToBuildAvroSchema(
+      Type type, Map<Type, Schema> cache, String recordName, String namespace) 
{
     switch (type.typeId()) {
       case RECORD:
         Types.RecordType record = (Types.RecordType) type;
         List<Schema> schemas = new ArrayList<>();
         record.fields().forEach(f -> {
-          Schema tempSchema = visitInternalSchemaToBuildAvroSchema(f.type(), 
cache, recordName + "_" + f.name());
+          Schema tempSchema = visitInternalSchemaToBuildAvroSchema(
+              f.type(), cache, recordName + "_" + f.name(), namespace);
           // convert tempSchema
           Schema result = f.isOptional() ? 
AvroInternalSchemaConverter.nullableSchema(tempSchema) : tempSchema;
           schemas.add(result);
@@ -281,13 +287,13 @@ public class AvroInternalSchemaConverter {
         if (recordSchema != null) {
           return recordSchema;
         }
-        recordSchema = visitInternalRecordToBuildAvroRecord(record, schemas, 
recordName);
+        recordSchema = visitInternalRecordToBuildAvroRecord(record, schemas, 
recordName, namespace);
         cache.put(record, recordSchema);
         return recordSchema;
       case ARRAY:
         Types.ArrayType array = (Types.ArrayType) type;
         Schema elementSchema;
-        elementSchema = 
visitInternalSchemaToBuildAvroSchema(array.elementType(), cache, recordName);
+        elementSchema = 
visitInternalSchemaToBuildAvroSchema(array.elementType(), cache, recordName, 
namespace);
         Schema arraySchema;
         arraySchema = cache.get(array);
         if (arraySchema != null) {
@@ -300,8 +306,8 @@ public class AvroInternalSchemaConverter {
         Types.MapType map = (Types.MapType) type;
         Schema keySchema;
         Schema valueSchema;
-        keySchema = visitInternalSchemaToBuildAvroSchema(map.keyType(), cache, 
recordName);
-        valueSchema = visitInternalSchemaToBuildAvroSchema(map.valueType(), 
cache, recordName);
+        keySchema = visitInternalSchemaToBuildAvroSchema(map.keyType(), cache, 
recordName, namespace);
+        valueSchema = visitInternalSchemaToBuildAvroSchema(map.valueType(), 
cache, recordName, namespace);
         Schema mapSchema;
         mapSchema = cache.get(map);
         if (mapSchema != null) {
@@ -321,7 +327,8 @@ public class AvroInternalSchemaConverter {
    * Converts hudi RecordType to Avro RecordType.
    * this is auxiliary function used by visitInternalSchemaToBuildAvroSchema
    */
-  private static Schema visitInternalRecordToBuildAvroRecord(Types.RecordType 
record, List<Schema> fieldSchemas, String recordName) {
+  private static Schema visitInternalRecordToBuildAvroRecord(
+      Types.RecordType record, List<Schema> fieldSchemas, String recordName, 
String namespace) {
     List<Types.Field> fields = record.fields();
     List<Schema.Field> avroFields = new ArrayList<>();
     for (int i = 0; i < fields.size(); i++) {
@@ -329,7 +336,7 @@ public class AvroInternalSchemaConverter {
       Schema.Field field = new Schema.Field(f.name(), fieldSchemas.get(i), 
f.doc(), f.isOptional() ? JsonProperties.NULL_VALUE : null);
       avroFields.add(field);
     }
-    return Schema.createRecord(recordName, null, null, false, avroFields);
+    return Schema.createRecord(recordName, null, namespace, false, avroFields);
   }
 
   /**
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/internal/schema/utils/AvroSchemaEvolutionUtils.java
 
b/hudi-common/src/main/java/org/apache/hudi/internal/schema/utils/AvroSchemaEvolutionUtils.java
index 520a6b9ec7..413a3c4df1 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/internal/schema/utils/AvroSchemaEvolutionUtils.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/internal/schema/utils/AvroSchemaEvolutionUtils.java
@@ -129,7 +129,9 @@ public class AvroSchemaEvolutionUtils {
     // try to correct all changes
     TableChanges.ColumnUpdateChange updateChange = 
TableChanges.ColumnUpdateChange.get(writeInternalSchema);
     candidateUpdateCols.stream().forEach(f -> 
updateChange.updateColumnNullability(f, true));
-    Schema result = 
AvroInternalSchemaConverter.convert(SchemaChangeUtils.applyTableChanges2Schema(writeInternalSchema,
 updateChange), writeSchema.getName());
+    Schema result = AvroInternalSchemaConverter.convert(
+        SchemaChangeUtils.applyTableChanges2Schema(writeInternalSchema, 
updateChange),
+        writeSchema.getName(), writeSchema.getNamespace());
     return result;
   }
 }

Reply via email to