yumeng 2 هفته پیش
والد
کامیت
60bcf4d0fb
1فایلهای تغییر یافته به همراه120 افزوده شده و 2 حذف شده
  1. 120 2
      src/main/java/com/adx/tencent/report/AdBidReportStore.java

+ 120 - 2
src/main/java/com/adx/tencent/report/AdBidReportStore.java

@@ -10,8 +10,10 @@ import java.time.LocalDate;
 import java.time.YearMonth;
 import java.time.ZoneId;
 import java.util.ArrayList;
+import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Locale;
+import java.util.Map;
 
 public class AdBidReportStore {
 
@@ -118,6 +120,7 @@ public class AdBidReportStore {
         try (PreparedStatement ps = conn.prepareStatement(sql)) {
             ps.execute();
         }
+        syncMissingColumns(conn, mainSchema, media.hotTable, reportSchema, table);
     }
 
     private void insertHistoryBatch(Connection conn,
@@ -125,8 +128,9 @@ public class AdBidReportStore {
                                     YearMonth month,
                                     List<Long> candidateIds) throws Exception {
         String table = historyTable(media, month);
-        String sql = "INSERT IGNORE INTO `" + reportSchema + "`.`" + table + "` " +
-                "SELECT * FROM `" + mainSchema + "`.`" + media.hotTable + "` " +
+        String columns = commonColumnList(conn, mainSchema, media.hotTable, reportSchema, table);
+        String sql = "INSERT IGNORE INTO `" + reportSchema + "`.`" + table + "` (" + columns + ") " +
+                "SELECT " + columns + " FROM `" + mainSchema + "`.`" + media.hotTable + "` " +
                 "WHERE id IN (" + placeholders(candidateIds.size()) + ") " +
                 "ORDER BY created_at ASC, id ASC";
         try (PreparedStatement ps = conn.prepareStatement(sql)) {
@@ -326,6 +330,112 @@ public class AdBidReportStore {
         }
     }
 
+    private void syncMissingColumns(Connection conn,
+                                    String sourceSchema,
+                                    String sourceTable,
+                                    String targetSchema,
+                                    String targetTable) throws Exception {
+        Map<String, ColumnDef> sourceColumns = tableColumns(conn, sourceSchema, sourceTable);
+        Map<String, ColumnDef> targetColumns = tableColumns(conn, targetSchema, targetTable);
+        String previousColumn = null;
+        for (ColumnDef column : sourceColumns.values()) {
+            if (!targetColumns.containsKey(column.name)) {
+                String sql = "ALTER TABLE `" + targetSchema + "`.`" + targetTable + "` ADD COLUMN " +
+                        columnDefinition(column) + afterClause(previousColumn);
+                try (PreparedStatement ps = conn.prepareStatement(sql)) {
+                    ps.execute();
+                }
+            }
+            previousColumn = column.name;
+        }
+    }
+
+    private String commonColumnList(Connection conn,
+                                    String sourceSchema,
+                                    String sourceTable,
+                                    String targetSchema,
+                                    String targetTable) throws Exception {
+        Map<String, ColumnDef> sourceColumns = tableColumns(conn, sourceSchema, sourceTable);
+        Map<String, ColumnDef> targetColumns = tableColumns(conn, targetSchema, targetTable);
+        List<String> columns = new ArrayList<>();
+        for (String column : sourceColumns.keySet()) {
+            if (targetColumns.containsKey(column)) {
+                columns.add("`" + column + "`");
+            }
+        }
+        if (columns.isEmpty()) {
+            throw new IllegalStateException("no common columns between " + sourceTable + " and " + targetTable);
+        }
+        return String.join(", ", columns);
+    }
+
+    private Map<String, ColumnDef> tableColumns(Connection conn, String schema, String table) throws Exception {
+        String sql = "SHOW FULL COLUMNS FROM `" + schema + "`.`" + table + "`";
+        Map<String, ColumnDef> columns = new LinkedHashMap<>();
+        try (PreparedStatement ps = conn.prepareStatement(sql);
+             ResultSet rs = ps.executeQuery()) {
+            while (rs.next()) {
+                ColumnDef column = new ColumnDef(
+                        rs.getString("Field"),
+                        rs.getString("Type"),
+                        rs.getString("Null"),
+                        rs.getString("Default"),
+                        rs.wasNull(),
+                        rs.getString("Extra"),
+                        rs.getString("Comment")
+                );
+                columns.put(column.name, column);
+            }
+        }
+        return columns;
+    }
+
+    private static String columnDefinition(ColumnDef column) {
+        StringBuilder ddl = new StringBuilder();
+        ddl.append("`").append(column.name).append("` ").append(column.type);
+        if ("NO".equalsIgnoreCase(column.nullable)) {
+            ddl.append(" NOT NULL");
+        } else {
+            ddl.append(" NULL");
+        }
+        if (!column.defaultWasNull) {
+            ddl.append(" DEFAULT ");
+            if (isSqlExpressionDefault(column.defaultValue)) {
+                ddl.append(column.defaultValue);
+            } else {
+                ddl.append("'").append(escapeSql(column.defaultValue)).append("'");
+            }
+        }
+        if (column.extra != null && !column.extra.isBlank()) {
+            ddl.append(" ").append(column.extra);
+        }
+        if (column.comment != null && !column.comment.isBlank()) {
+            ddl.append(" COMMENT '").append(escapeSql(column.comment)).append("'");
+        }
+        return ddl.toString();
+    }
+
+    private static String afterClause(String previousColumn) {
+        return previousColumn == null ? " FIRST" : " AFTER `" + previousColumn + "`";
+    }
+
+    private static boolean isSqlExpressionDefault(String value) {
+        if (value == null) {
+            return false;
+        }
+        String upper = value.toUpperCase(Locale.ROOT);
+        return "CURRENT_TIMESTAMP".equals(upper)
+                || upper.startsWith("CURRENT_TIMESTAMP(")
+                || "NULL".equals(upper)
+                || "TRUE".equals(upper)
+                || "FALSE".equals(upper)
+                || upper.matches("-?\\d+(\\.\\d+)?");
+    }
+
+    private static String escapeSql(String value) {
+        return value == null ? "" : value.replace("\\", "\\\\").replace("'", "''");
+    }
+
     private int deleteHotBatch(Connection conn,
                                String hotTable,
                                String historyTable,
@@ -379,6 +489,14 @@ public class AdBidReportStore {
         return value;
     }
 
+    private record ColumnDef(String name,
+                             String type,
+                             String nullable,
+                             String defaultValue,
+                             boolean defaultWasNull,
+                             String extra,
+                             String comment) {}
+
     public static class RunSummary {
         public long archivedRows;
         public long archivedTrackingRows;