Упростите работу ETL, напишите слой клея Canal

задняя часть

помещение

Это статья, которую я долго сдерживал, всегда хотел ее написать, но все время забывал написать. Вся статья может быть немного запущенной, и в ней относительно подробно описывается, как написать небольшой «фреймворк». Этот мощный связующий слой служит в производственной среде более полугода, здесь мы пытаемся удалить код спаренного бизнеса и извлечь относительно сжатую версию.

В одной из статей, которые я написал ранее, упоминалосьCanalРазобратьMySQLизbinlogОбъект после события выглядит следующим образом (получено изCanalисходный кодcom.alibaba.otter.canal.protocol.FlatMessage):

Если исходный объект анализируется напрямую, будет много кода шаблона синтаксического анализа, и как только произойдет изменение, будет затронуто все тело, чего мы не хотим. Поэтому я потратил немного времени на написаниеCanalклеевой слой, позволяющий получитьFlatMessageПреобразование непосредственно в соответствующее имя таблицыDTOНапример, это может в определенной степени повысить эффективность разработки и уменьшить шаблонный код.Схема потока данных этого связующего слоя выглядит следующим образом:

Для написания такого клеевого слоя в основном используют:

  • отражение.
  • аннотация.
  • режим стратегии.
  • IOCконтейнер (по желанию).

Модули проекта следующие:

  • canal-glue-core: Основные функции.
  • spring-boot-starter-canal-glue:приспособлениеSpringизIOCконтейнер, добавьте автоконфигурацию.
  • canal-glue-example: Используйте примеры и тесты.

В следующем разделе будет подробно проанализировано, как реализован этот клеевой слой.

импортировать зависимости

Чтобы не загрязнять внешние сервисные зависимости, ссылающиеся на этот модуль, за исключениемJSONВ дополнение к преобразованным зависимостям, другие зависимостиscopeопределяется какprovideилиtestтип, версия зависимостей иBOMследующим образом:

<properties>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
        <maven.compiler.source>1.8</maven.compiler.source>
        <maven.compiler.target>1.8</maven.compiler.target>
        <spring.boot.version>2.3.0.RELEASE</spring.boot.version>
        <maven.compiler.plugin.version>3.8.1</maven.compiler.plugin.version>
        <lombok.version>1.18.12</lombok.version>
        <fastjson.version>1.2.73</fastjson.version>
</properties>
<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-dependencies</artifactId>
            <version>${spring.boot.version}</version>
            <scope>import</scope>
            <type>pom</type>
        </dependency>
    </dependencies>
</dependencyManagement>
<dependencies>
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <version>${lombok.version}</version>
        <scope>provided</scope>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-test</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter</artifactId>
        <scope>provided</scope>
    </dependency>
    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>fastjson</artifactId>
        <version>${fastjson.version}</version>
    </dependency>
</dependencies>

в,canal-glue-coreМодули по существу зависят только отfastjson, можно полностью отключитьspringиспользование системы.

Базовая архитектура

Вот архитектурная схема "задним числом", потому что для того, чтобы быстро выйти в онлайн, первая версия не так уж много учитывала, да еще и бизнес-код спаривала, а компоненты извлекались позже:

Модуль конфигурации дизайна (удален)

Модуль конфигурации дизайна учитывает использование внешних файлов конфигурации и чистых аннотаций при проектировании. На раннем этапе использовался внешний файл конфигурации JSON, а чистые аннотации были добавлены позже. Выберите один из двух. В этом разделе кратко представлена ​​загрузка конфигурации внешнего файла конфигурации JSON, а чистые аннотации оставлены для анализа процессорного модуля позже.

Сначала я хотел быстро разработать клеевой слой, поэтому в файле конфигурации использовался более читаемый файл конфигурации.JSONФормат:

{
  "version": 1,
  "module": "canal-glue",
  "databases": [
    {
      "database": "db_payment_service",
      "processors": [
        {
          "table": "payment_order",
          "processor": "x.y.z.PaymentOrderProcessor",
          "exceptionHandler": "x.y.z.PaymentOrderExceptionHandler"
        }
      ]
    },
    {
      ......
    }
  ]
}

При разработке конфигурации JSON старайтесь не использовать JSON Array в качестве конфигурации верхнего уровня, потому что объект, спроектированный таким образом, будет странным.

Поскольку приложения, использующие этот модуль, могут иметь дело сCanalсинтаксический анализ нескольких баз данных вышестоящего уровняbinlogсобытий, поэтому при настройке дизайна модуля необходимо использоватьdatabaseзаKEY, смонтировать несколькоtableи соответствующая таблицаbinlogОбработчики событий и обработчики исключений. затем лицом к лицуJSONФормат файла развертывает соответствующий класс сущности:

@Data
public class CanalGlueProcessorConf {

    private String table;

    private String processor;

    private String exceptionHandler;
}

@Data
public class CanalGlueDatabaseConf {

    private String database;

    private List<CanalGlueProcessorConf> processors;
}

@Data
public class CanalGlueConf {

    private Long version;

    private String module;

    private List<CanalGlueDatabaseConf> database;
}

После написания сущности можно написать загрузчик конфигурации.Для простоты файл конфигурации можно разместить непосредственноClassPathНиже загрузчик выглядит следующим образом:

public interface CanalGlueConfLoader {

    CanalGlueConf load(String location);
}

// 实现
public class ClassPathCanalGlueConfLoader implements CanalGlueConfLoader {

    @Override
    public CanalGlueConf load(String location) {
        ClassPathResource resource = new ClassPathResource(location);
        Assert.isTrue(resource.exists(), String.format("类路径下不存在文件%s", location));
        try {
            String content = StreamUtils.copyToString(resource.getInputStream(), StandardCharsets.UTF_8);
            return JSON.parseObject(content, CanalGlueConf.class);
        } catch (IOException e) {
            // should not reach
            throw new IllegalStateException(e);
        }
    }
}

читатьClassPathодин из следующихlocationстрока содержимого файла для абсолютного пути, затем используйтеFasfjsonПревратиться вCanalGlueConfобъект. Это реализация по умолчанию, использующаяcanal-glueМодули могут переопределить эту реализацию, чтобы загрузить конфигурацию с пользовательской реализацией.

От модуля конфигурации JSON полностью отказались, когда позже из бизнес-системы был извлечен связующий слой, и он был реализован с использованием чистого наследования компонентов, основанных на аннотациях и базовой абстракции.

Разработка основного модуля

В основном включает несколько модулей:

  • Определение базовой модели.
  • Разработка слоя адаптера.
  • Разработка уровня преобразователя и парсера.
  • Развитие уровня процессора.
  • Разработка модуля автоконфигурации глобального компонента (ограниченоSpringсистема, была извлечена вspring-boot-starter-canal-glueмодуль).
  • CanalGlueразработка.

Определение базовой модели

определить верхний уровеньKEY, то есть для определенной таблицы в БД требуется уникальный идентификатор:

// 模型表对象
public interface ModelTable {

    String database();

    String table();

    static ModelTable of(String database, String table) {
        return DefaultModelTable.of(database, table);
    }
}

@RequiredArgsConstructor(access = AccessLevel.PACKAGE, staticName = "of")
public class DefaultModelTable implements ModelTable {

    private final String database;
    private final String table;

    @Override
    public String database() {
        return database;
    }

    @Override
    public String table() {
        return table;
    }

    @Override
    public boolean equals(Object o) {
        if (this == o) {
            return true;
        }
        if (o == null || getClass() != o.getClass()) {
            return false;
        }
        DefaultModelTable that = (DefaultModelTable) o;
        return Objects.equals(database, that.database) &&
                Objects.equals(table, that.table);
    }

    @Override
    public int hashCode() {
        return Objects.hash(database, table);
    }
}

Класс реализации здесьDefaultModelTableпереписанныйequals()а такжеhashCode()способ облегчитьModelTableПример приложенияHashMapконтейнерKEY, чтобы его можно было спроектировать позжеModelTable -> Processorструктура кэша.

потому чтоCanalслужитьKafkaСодержимое события представляет собой необработанную строку, поэтому определитеFlatMessageВ основном тот же класс событийCanalBinLogEvent:

@Data
public class CanalBinLogEvent {

    /**
     * 事件ID,没有实际意义
     */
    private Long id;

    /**
     * 当前更变后节点数据
     */
    private List<Map<String, String>> data;

    /**
     * 主键列名称列表
     */
    private List<String> pkNames;

    /**
     * 当前更变前节点数据
     */
    private List<Map<String, String>> old;

    /**
     * 类型 UPDATE\INSERT\DELETE\QUERY
     */
    private String type;

    /**
     * binlog execute time
     */
    private Long es;

    /**
     * dml build timestamp
     */
    private Long ts;

    /**
     * 执行的sql,不一定存在
     */
    private String sql;

    /**
     * 数据库名称
     */
    private String database;

    /**
     * 表名称
     */
    private String table;

    /**
     * SQL类型映射
     */
    private Map<String, Integer> sqlType;

    /**
     * MySQL字段类型映射
     */
    private Map<String, String> mysqlType;

    /**
     * 是否DDL
     */
    private Boolean isDdl;
}

В соответствии с этим объектом события определите объект результата после синтаксического анализаCanalBinLogResult:

// 常量
@RequiredArgsConstructor
@Getter
public enum BinLogEventType {
    
    QUERY("QUERY", "查询"),

    INSERT("INSERT", "新增"),

    UPDATE("UPDATE", "更新"),

    DELETE("DELETE", "删除"),

    ALTER("ALTER", "列修改操作"),

    UNKNOWN("UNKNOWN", "未知"),

    ;

    private final String type;
    private final String description;

    public static BinLogEventType fromType(String type) {
        for (BinLogEventType binLogType : BinLogEventType.values()) {
            if (binLogType.getType().equals(type)) {
                return binLogType;
            }
        }
        return BinLogEventType.UNKNOWN;
    }
}

// 常量
@RequiredArgsConstructor
@Getter
public enum OperationType {

    /**
     * DML
     */
    DML("dml", "DML语句"),

    /**
     * DDL
     */
    DDL("ddl", "DDL语句"),
    ;

    private final String type;
    private final String description;
}

@Data
public class CanalBinLogResult<T> {

    /**
     * 提取的长整型主键
     */
    private Long primaryKey;


    /**
     * binlog事件类型
     */
    private BinLogEventType binLogEventType;

    /**
     * 更变前的数据
     */
    private T beforeData;

    /**
     * 更变后的数据
     */
    private T afterData;

    /**
     * 数据库名称
     */
    private String databaseName;

    /**
     * 表名称
     */
    private String tableName;

    /**
     * sql语句 - 一般是DDL的时候有用
     */
    private String sql;

    /**
     * MySQL操作类型
     */
    private OperationType operationType;
}

Уровень адаптера разработки

Определите адаптер верхнего уровняSPIинтерфейс:

public interface SourceAdapter<SOURCE, SINK> {

    SINK adapt(SOURCE source);
}

Затем разработайте класс реализации адаптера:

// 原始字符串直接返回
@RequiredArgsConstructor(access = AccessLevel.PACKAGE, staticName = "of")
class RawStringSourceAdapter implements SourceAdapter<String, String> {

    @Override
    public String adapt(String source) {
        return source;
    }
}

// Fastjson转换
@RequiredArgsConstructor(access = AccessLevel.PACKAGE, staticName = "of")
class FastJsonSourceAdapter<T> implements SourceAdapter<String, T> {

    private final Class<T> klass;

    @Override
    public T adapt(String source) {
        if (StringUtils.isEmpty(source)) {
            return null;
        }
        return JSON.parseObject(source, klass);
    }
}

// Facade
public enum SourceAdapterFacade {

    /**
     * 单例
     */
    X;

    private static final SourceAdapter<String, String> I_S_A = RawStringSourceAdapter.of();

    @SuppressWarnings("unchecked")
    public <T> T adapt(Class<T> klass, String source) {
        if (klass.isAssignableFrom(String.class)) {
            return (T) I_S_A.adapt(source);
        }
        return FastJsonSourceAdapter.of(klass).adapt(source);
    }
}

наконец, используется напрямуюSourceAdapterFacade#adapt()метод, потому что на практике только необработанные строки иString -> Class实例, дизайн слоя адаптера может быть проще.

Разработка слоев преобразователя и парсера

заCanalразбор завершенbinlogсобытие,dataа такжеoldсобственностьK-Vструктура иKEYобаStringТип, который требует обходного синтаксического анализа для вывода полного целевого экземпляра.

Тип атрибута преобразованного экземпляра в настоящее время поддерживает только классы-оболочки, а примитивные типы, такие как int, — нет.

Введены две пользовательские аннотации для лучшего сопоставления между целевыми объектами и фактической базой данных, именами таблиц и столбцов, типами столбцов.CanalModelа также@CanalField, они определяются следующим образом:

// @CanalModel
@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.TYPE)
public @interface CanalModel {

    /**
     * 目标数据库
     */
    String database();

    /**
     * 目标表
     */
    String table();

    /**
     * 属性名 -> 列名命名转换策略,可选值有:DEFAULT(原始)、UPPER_UNDERSCORE(驼峰转下划线大写)和LOWER_UNDERSCORE(驼峰转下划线小写)
     */
    FieldNamingPolicy fieldNamingPolicy() default FieldNamingPolicy.DEFAULT;
}

// @CanalField
@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.FIELD)
public @interface CanalField {

    /**
     * 行名称
     *
     * @return columnName
     */
    String columnName() default "";

    /**
     * sql字段类型
     *
     * @return JDBCType
     */
    JDBCType sqlType() default JDBCType.NULL;

    /**
     * 转换器类型
     *
     * @return klass
     */
    Class<? extends BaseCanalFieldConverter<?>> converterKlass() default NullCanalFieldConverter.class;
}

Определить интерфейс преобразователя верхнего уровняBinLogFieldConverter:

public interface BinLogFieldConverter<SOURCE, TARGET> {

    TARGET convert(SOURCE source);
}

В настоящее время он предварительно доступен через целевой атрибутClassи указан через аннотациюSQLTypeТип для соответствия, поэтому определите другой абстрактный преобразовательBaseCanalFieldConverter:

public abstract class BaseCanalFieldConverter<T> implements BinLogFieldConverter<String, T> {

    private final SQLType sqlType;
    private final Class<?> klass;

    protected BaseCanalFieldConverter(SQLType sqlType, Class<?> klass) {
        this.sqlType = sqlType;
        this.klass = klass;
    }

    @Override
    public T convert(String source) {
        if (StringUtils.isEmpty(source)) {
            return null;
        }
        return convertInternal(source);
    }

    /**
     * 内部转换方法
     *
     * @param source 源字符串
     * @return T
     */
    protected abstract T convertInternal(String source);

    /**
     * 返回SQL类型
     *
     * @return SQLType
     */
    public SQLType sqlType() {
        return sqlType;
    }

    /**
     * 返回类型
     *
     * @return Class<?>
     */
    public Class<?> typeKlass() {
        return klass;
    }
}

BaseCanalFieldConverterориентирован на один атрибут в целевом экземпляре, например.LongСвойство типа, которое реализуетBigIntCanalFieldConverter:

public class BigIntCanalFieldConverter extends BaseCanalFieldConverter<Long> {

    /**
     * 单例
     */
    public static final BaseCanalFieldConverter<Long> X = new BigIntCanalFieldConverter();

    private BigIntCanalFieldConverter() {
        super(JDBCType.BIGINT, Long.class);
    }

    @Override
    protected Long convertInternal(String source) {
        if (null == source) {
            return null;
        }
        return Long.valueOf(source);
    }
}

Другие типы и так далее, наиболее часто используемые встроенные преобразователи, которые были разработаны, следующие:

JDBCType JAVAType преобразователь
NULL Void NullCanalFieldConverter
BIGINT Long BigIntCanalFieldConverter
VARCHAR String VarcharCanalFieldConverter
DECIMAL BigDecimal DecimalCanalFieldConverter
INTEGER Integer IntCanalFieldConverter
TINYINT Integer TinyIntCanalFieldConverter
DATE java.time.LocalDate SqlDateCanalFieldConverter0
DATE java.sql.Date SqlDateCanalFieldConverter1
TIMESTAMP java.time.LocalDateTime TimestampCanalFieldConverter0
TIMESTAMP java.util.Date TimestampCanalFieldConverter1
TIMESTAMP java.time.OffsetDateTime TimestampCanalFieldConverter2

Все реализации преобразователя разработаны как синглтоны без сохранения состояния для простой динамической регистрации и переопределения. Затем определите фабрику конвертеровCanalFieldConverterFactory,поставкаAPIЗагрузите экземпляр целевого преобразователя, указав параметры:

// 入参
@SuppressWarnings("rawtypes")
@Builder
@Data
public class CanalFieldConvertInput {

    private Class<?> fieldKlass;
    private Class<? extends BaseCanalFieldConverter> converterKlass;
    private SQLType sqlType;

    @Tolerate
    public CanalFieldConvertInput() {

    }
}

// 结果
@Builder
@Getter
public class CanalFieldConvertResult {

    private final BaseCanalFieldConverter<?> converter;
}

// 接口
public interface CanalFieldConverterFactory {

    default void registerConverter(BaseCanalFieldConverter<?> converter) {
        registerConverter(converter, true);
    }

    void registerConverter(BaseCanalFieldConverter<?> converter, boolean replace);

    CanalFieldConvertResult load(CanalFieldConvertInput input);
}

CanalFieldConverterFactoryПредоставляет возможность регистрации пользовательских конвертеровregisterConverter()метод, который позволяет пользователям регистрировать пользовательские преобразователи и переопределять преобразователи по умолчанию.

На этом этапе вы можете загрузить преобразователь атрибута экземпляра с помощью указанных параметров, получить экземпляр преобразователя, а затем проанализировать соответствующий целевой экземпляр из исходного события.K-Vструктура. Затем вам нужно написать основной модуль парсера, который в основном включает в себя три аспекта:

  • ТолькоBIGINTВведите разрешение первичного ключа (это железное правило технического задания компании,MySQLКаждая таблица может определять только уникальныеBIGINT UNSIGNEDПервичный ключ тенденции автоинкремента).
  • Данные до изменения, соответствующие данным в исходном событииoldузел атрибута (не обязательно присутствует, напримерINSERTЭтот узел свойства не существует в операторе).
  • Измененные данные, соответствующие исходному событиюdataузел свойства.

Определить интерфейс парсераCanalBinLogEventParserследующим образом:

public interface CanalBinLogEventParser {

    /**
     * 解析binlog事件
     *
     * @param event               事件
     * @param klass               目标类型
     * @param primaryKeyFunction  主键映射方法
     * @param commonEntryFunction 其他属性映射方法
     * @return CanalBinLogResult
     */
    <T> List<CanalBinLogResult<T>> parse(CanalBinLogEvent event,
                                         Class<T> klass,
                                         BasePrimaryKeyTupleFunction primaryKeyFunction,
                                         BaseCommonEntryFunction<T> commonEntryFunction);
}

Метод разбора парсера зависит от:

  • binlogЭкземпляр события, это результат работы вышестоящего компонента адаптера.
  • Целевой тип конверсии.
  • BasePrimaryKeyTupleFunctionЭкземпляр метода сопоставления первичного ключа, по умолчанию используется встроенныйBigIntPrimaryKeyTupleFunction.
  • BaseCommonEntryFunctionЭкземпляр общего метода сопоставления свойств столбца, отличного от первичного ключа, встроенный используется по умолчаниюReflectionBinLogEntryFunction(Это ядро ​​преобразования столбцов, не являющихся первичными ключами, которое использует отражение.).

Результат разбора -List, Причина вFlatMessageСтруктура данных при пакетной записи изначальноList<Map<String,String>>, вот просто "толкай лодку".

уровень процессора разработки

Процессор — это точка входа разработчика для обработки окончательно разобранной сущности, ему нужно только выбрать соответствующий метод обработки для разных типов событий, который выглядит следующим образом:

public abstract class BaseCanalBinlogEventProcessor<T> extends BaseParameterizedTypeReferenceSupport<T> {

    protected void processInsertInternal(CanalBinLogResult<T> result) {
    }

    protected void processUpdateInternal(CanalBinLogResult<T> result) {
    }

    protected void processDeleteInternal(CanalBinLogResult<T> result) {
    }

    protected void processDDLInternal(CanalBinLogResult<T> result) {
    }
}

Например, нужно обработатьInsertсобытие, подкласс наследуетBaseCanalBinlogEventProcessor, соответствующий класс сущностей (замена дженериков) использует@CanalModelАннотируйте объявление, затем переопределитеprocessInsertInternal()метод. Подпроцессоры периода могут переопределять пользовательские экземпляры обработчиков исключений, например:

@Override
protected ExceptionHandler exceptionHandler() {
    return EXCEPTION_HANDLER;
}

/**
    * 覆盖默认的ExceptionHandler.NO_OP
    */
private static final ExceptionHandler EXCEPTION_HANDLER = (event, throwable)
        -> log.error("解析binlog事件出现异常,事件内容:{}", JSON.toJSONString(event), throwable);

Кроме того, в некоторых сценариях необходимо специализировать результаты до или после обратного вызова, поэтому вводится реализация перехватчика (цепочки) результатов синтаксического анализа и соответствующий класс.BaseParseResultInterceptor:

public abstract class BaseParseResultInterceptor<T> extends BaseParameterizedTypeReferenceSupport<T> {

    public BaseParseResultInterceptor() {
        super();
    }

    public void onParse(ModelTable modelTable) {

    }

    public void onBeforeInsertProcess(ModelTable modelTable, T beforeData, T afterData) {

    }

    public void onAfterInsertProcess(ModelTable modelTable, T beforeData, T afterData) {

    }

    public void onBeforeUpdateProcess(ModelTable modelTable, T beforeData, T afterData) {

    }

    public void onAfterUpdateProcess(ModelTable modelTable, T beforeData, T afterData) {

    }

    public void onBeforeDeleteProcess(ModelTable modelTable, T beforeData, T afterData) {

    }

    public void onAfterDeleteProcess(ModelTable modelTable, T beforeData, T afterData) {

    }

    public void onBeforeDDLProcess(ModelTable modelTable, T beforeData, T afterData, String sql) {

    }

    public void onAfterDDLProcess(ModelTable modelTable, T beforeData, T afterData, String sql) {

    }

    public void onParseFinish(ModelTable modelTable) {

    }

    public void onParseCompletion(ModelTable modelTable) {

    }
}

Время обратного вызова перехватчика результатов синтаксического анализа см. на приведенной выше архитектурной диаграмме илиBaseCanalBinlogEventProcessorисходный код.

Разработать модуль автоконфигурации глобального компонента

если используетсяSpringКонтейнер, вам нужно добавить класс конфигурации для загрузки всех существующих компонентов, добавить глобальный класс конфигурацииCanalGlueAutoConfiguration(Этот класс можно найти в папке проектаspring-boot-starter-canal-glueКак вы можете видеть в модуле, этот модуль имеет только один класс):

@Configuration
public class CanalGlueAutoConfiguration implements SmartInitializingSingleton, BeanFactoryAware {

    private ConfigurableListableBeanFactory configurableListableBeanFactory;

    @Bean
    @ConditionalOnMissingBean
    public CanalBinlogEventProcessorFactory canalBinlogEventProcessorFactory() {
        return InMemoryCanalBinlogEventProcessorFactory.of();
    }

    @Bean
    @ConditionalOnMissingBean
    public ModelTableMetadataManager modelTableMetadataManager(CanalFieldConverterFactory canalFieldConverterFactory) {
        return InMemoryModelTableMetadataManager.of(canalFieldConverterFactory);
    }

    @Bean
    @ConditionalOnMissingBean
    public CanalFieldConverterFactory canalFieldConverterFactory() {
        return InMemoryCanalFieldConverterFactory.of();
    }

    @Bean
    @ConditionalOnMissingBean
    public CanalBinLogEventParser canalBinLogEventParser() {
        return DefaultCanalBinLogEventParser.of();
    }

    @Bean
    @ConditionalOnMissingBean
    public ParseResultInterceptorManager parseResultInterceptorManager(ModelTableMetadataManager modelTableMetadataManager) {
        return InMemoryParseResultInterceptorManager.of(modelTableMetadataManager);
    }

    @Bean
    @Primary
    public CanalGlue canalGlue(CanalBinlogEventProcessorFactory canalBinlogEventProcessorFactory) {
        return DefaultCanalGlue.of(canalBinlogEventProcessorFactory);
    }

    @Override
    public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
        this.configurableListableBeanFactory = (ConfigurableListableBeanFactory) beanFactory;
    }

    @SuppressWarnings({"rawtypes", "unchecked"})
    @Override
    public void afterSingletonsInstantiated() {
        ParseResultInterceptorManager parseResultInterceptorManager
                = configurableListableBeanFactory.getBean(ParseResultInterceptorManager.class);
        ModelTableMetadataManager modelTableMetadataManager
                = configurableListableBeanFactory.getBean(ModelTableMetadataManager.class);
        CanalBinlogEventProcessorFactory canalBinlogEventProcessorFactory
                = configurableListableBeanFactory.getBean(CanalBinlogEventProcessorFactory.class);
        CanalBinLogEventParser canalBinLogEventParser
                = configurableListableBeanFactory.getBean(CanalBinLogEventParser.class);
        Map<String, BaseParseResultInterceptor> interceptors
                = configurableListableBeanFactory.getBeansOfType(BaseParseResultInterceptor.class);
        interceptors.forEach((k, interceptor) -> parseResultInterceptorManager.registerParseResultInterceptor(interceptor));
        Map<String, BaseCanalBinlogEventProcessor> processors
                = configurableListableBeanFactory.getBeansOfType(BaseCanalBinlogEventProcessor.class);
        processors.forEach((k, processor) -> processor.init(canalBinLogEventParser, modelTableMetadataManager,
                canalBinlogEventProcessorFactory, parseResultInterceptorManager));
    }
}

Чтобы другие службы могли лучше вводить этот класс конфигурации, вы можете использоватьspring.factoriesхарактеристики. новыйresources/META-INF/spring.factoriesфайл со следующим содержимым:

org.springframework.boot.autoconfigure.EnableAutoConfiguration=cn.throwx.canal.gule.config.CanalGlueAutoConfiguration

Таким образом, вводяspring-boot-starter-canal-glueВы можете активировать все используемые компоненты и инициализировать все компоненты, которые были добавлены вSpringПроцессор в контейнере.

CanalGlue Development

CanalGlueФактически он обеспечиваетbinlogЗапись обработки строки события, которая в настоящее время определена как интерфейс:

public interface CanalGlue {

    void process(String content);
}

реализация этого интерфейсаDefaultCanalGlueТоже очень просто:

@RequiredArgsConstructor(access = AccessLevel.PUBLIC, staticName = "of")
public class DefaultCanalGlue implements CanalGlue {

    private final CanalBinlogEventProcessorFactory canalBinlogEventProcessorFactory;

    @Override
    public void process(String content) {
        CanalBinLogEvent event = SourceAdapterFacade.X.adapt(CanalBinLogEvent.class, content);
        ModelTable modelTable = ModelTable.of(event.getDatabase(), event.getTable());
        canalBinlogEventProcessorFactory.get(modelTable).forEach(processor -> processor.process(event));
    }
}

Используйте исходный адаптер для преобразования строки вCanalBinLogEventэкземпляр, а затем поручить фабрике процессоров найти соответствующийBaseCanalBinlogEventProcessorСписок экземпляров событий для обработки ввода.

Используйте канальный клей

Он в основном включает в себя следующие измерения, все из которыхcanal-glue-exampleизtestПод пакетом:

  • Обычно обрабатывается процессоромINSERTмероприятие.
  • обычай дляDDLИзменен родительский обработчик предупреждений, реализованDDLОповещение об изменении.
  • Одна таблица соответствует нескольким процессорам.
  • Используйте обработчики результатов синтаксического анализа для определенных полейAESОбработка шифрования и дешифрования.
  • НетSpringПод контейнером он вообще используется программно.
  • использоватьopenjdk-jmhпровестиBenchmarkКонтрольные тесты производительности.

Здесь кратко упоминается, что вSpringИспользование под систему, введение зависимостейspring-boot-starter-canal-glue:

<dependency>
    <groupId>cn.throwx</groupId>
    <artifactId>spring-boot-starter-canal-glue</artifactId>
    <version>版本号</version>
</dependency>

написать сущность илиDTOДобрыйOrderModel:

@Data
@CanalModel(database = "db_order_service", table = "t_order", fieldNamingPolicy = FieldNamingPolicy.LOWER_UNDERSCORE)
public static class OrderModel {

    private Long id;

    private String orderId;

    private OffsetDateTime createTime;

    private BigDecimal amount;
}

используется здесь@CanalModelАннотация связывает базу данныхdb_order_serviceи столt_order, стратегия сопоставления имени атрибута с именем столбцаCamelCase в нижнее подчеркивание. Затем определите обработчикOrderProcessorи пользовательский обработчик исключений (необязательно, здесь для имитации создания пользовательских исключений при обработке событий):

@Component
public class OrderProcessor extends BaseCanalBinlogEventProcessor<OrderModel> {

    @Override
    protected void processInsertInternal(CanalBinLogResult<OrderModel> result) {
        OrderModel orderModel = result.getAfterData();
        logger.info("接收到订单保存binlog,主键:{},模拟抛出异常...", orderModel.getId());
        throw new RuntimeException(String.format("[id:%d]", orderModel.getId()));
    }

    @Override
    protected ExceptionHandler exceptionHandler() {
        return EXCEPTION_HANDLER;
    }

    /**
        * 覆盖默认的ExceptionHandler.NO_OP
        */
    private static final ExceptionHandler EXCEPTION_HANDLER = (event, throwable)
            -> log.error("解析binlog事件出现异常,事件内容:{}", JSON.toJSONString(event), throwable);
}

Предположим, данные о порядке записиbinlogСобытия следующие:

{
  "data": [
    {
      "id": "1",
      "order_id": "10086",
      "amount": "999.0",
      "create_time": "2020-03-02 05:12:49"
    }
  ],
  "database": "db_order_service",
  "es": 1583143969000,
  "id": 3,
  "isDdl": false,
  "mysqlType": {
    "id": "BIGINT",
    "order_id": "VARCHAR(64)",
    "amount": "DECIMAL(10,2)",
    "create_time": "DATETIME"
  },
  "old": null,
  "pkNames": [
    "id"
  ],
  "sql": "",
  "sqlType": {
    "id": -5,
    "order_id": 12,
    "amount": 3,
    "create_time": 93
  },
  "table": "t_order",
  "ts": 1583143969460,
  "type": "INSERT"
}

Результат выполнения следующий:

Если напрямую подключенCanalслужитьKafkaизTopicЭто также очень просто, сKafkaПример потребительского использования выглядит следующим образом:

@Slf4j
@Component
@RequiredArgsConstructor
public class CanalEventListeners {

    private final CanalGlue canalGlue;

    @KafkaListener(
            id = "${canal.event.order.listener.id:db-order-service-listener}",
            topics = "db_order_service", 
            containerFactory = "kafkaListenerContainerFactory"
    )
    public void onCrmMessage(String content) {
        canalGlue.process(content);
    }    
}

резюме

Автор разработал этоcanal-glueПервоначальное намерение состояло в том, чтобы сделать крупномасштабный преобразователь строк, который значительно повысит эффективность, потому что он только что вступил в контакт с полем «малых данных», а рабочая сила недостаточна, и ему необходимо иметь дело с большим количеством нисходящих потоков. отчеты, потому что невозможно тратить много рабочей силы на их непрерывную обработку повторяющимся шаблонным кодом. Хотя общий дизайн не очень элегантный,Хотя бы в плане повышения эффективности разработки,canal-glueВыполнено.

Репозиторий проекта:

  • Gitee:https://gitee.com/throwableDoge/canal-glue

Последний код склада временно размещен вdevelopфилиал.

(C-15-d e-a-20201005 висит уже почти месяц после окончания этой статьи)