【问题标题】:Delta load with Talend使用 Talend 进行增量负载
【发布时间】:2015-11-26 13:35:08
【问题描述】:

我是使用 Talend 的新手。 我想在我的 ETL 中使用增量负载。 我正在从 Mysql 数据源中提取并加载到 Postgresql 数据库中。 Mysql 数据源有 created_at 和 updated_at 时间戳,我想用它们来提取新数据或更新数据。 我之前已经在带有 SSIS 的 Sql Server 中实现了这一点。 我不确定如何使用 Talend 实施。 有人使用 Talend 实现了带时间戳的增量加载吗? 提前致谢。

【问题讨论】:

  • 你是如何在 sql server 中做到的?在 talend 中,一种方法是在 tMysqlInput 组件中编写 sql 查询并在此 sql 查询中添加 where 子句 - select * from sourcetable where created_dt >(您的最后加载日期)或 updated_dt>(您的最后加载日期)......你必须将上次加载日期存储在系统中并检索它以在此过滤器中使用。
  • 蛮力的方法是创建一个stage表,将最近几天的变化拉到stage中,然后从target中删除stage中存在的所有记录,然后追加整个stage达到目标。此外,您还必须创建另一个步骤来删除不存在的记录。只需将所有主键列从源拉到阶段,然后从目标中删除阶段中不存在的记录。 (通过主键加入)
  • @garpitmzn。我想我做了类似的事情。我有一个带有 created_at 和 updated_at 字段的查找表。在加载表格之前,如果数据是新的,我会检查该表格。我使用 SSIS 查找根据数据是新的还是更新的数据来拆分数据,并根据它插入或更新表。

标签: talend


【解决方案1】:

由于您已提交日期以识别增量,这将很容易。 有两个文件一个有当前日期(当工作流程开始时)。另一个文件的最后运行日期为低日期 19000101。在第一次加载中运行作业并从最后运行日期文件中读取日期,将此作为 where 子句来检查源数据时间戳 col>运行日期并运行作业。接下来将当前日期文件值移动到下一个运行日期文件。在增量中再次运行相同的过程。所以你会得到增量记录。

【讨论】:

    【解决方案2】:

    我为其中一个项目做过。 首先,我们将创建一个日志表,列如 Job_id、Job_name、start_time、end_time、status。 每当您运行一次作业时,您都需要在作业结束时更新此表。

    对于下一次运行,首先我们必须检查上次成功完成作业的时间,从上次作业开始时间开始,并将其放入一个变量中。

    在下面为表格输入提供条件

    created_at > 上次开始时间的变量

    见下图 或

    updated_at > 上次开始时间的变量

    check below image for job flow

    【讨论】:

      【解决方案3】:

      存在几个选项——这个答案并不特定于时间戳,但仍然会有所帮助。

      1. 您可以将内置 CDC 直接用于某些数据库(在 Talend 之外),例如 SQL Server 和 Oracle。可能与您的情况无关。

      2. 您可以将内置 CDC 用于某些数据库(在 Talend 内部)。需要订阅版本。包括 MySQL、Oracle、DB2、PostgreSQL、Sybase、MS SQL Server、Informix、Ingres 和 Teradata。 https://help.talend.com/reader/4UeRbZs9GU5n8b9nm3hUrQ/8yztvpROOkQauOWwo_0twA

      3. 您可以通过 SQL 和 Java(创建自定义例程)或包含的组件 tAddCDCRow 执行手动方法。这将涉及对键和/或非键值使用 MD5 哈希值。 * 始终对非关键字段使用哈希。 A. 组合键的散列,并包含比较的剩余字段。
        B. 在比较中包含所有关键字段并使用哈希值表示剩余。
        C. 对组合键使用一个哈希值,对其余字段使用另一个哈希值。

      如果使用 tAddCDCRow,则将一个组件用于组合键,将另一个组件用于其余字段。当然,如果只对键进行散列,则只需要一个组件。 如果使用自定义的 Java 函数,根据需要调用一次或两次。

      哈希Java函数:

      // data = Text for which to generate hash value.  Will be a combination of one or 
      more fields.  Recommend padding strings with spaces to accurately combine.  Can use |
      or similar dividers.  Convert integers and other non-text to strings, as necessary.
      public static String getMD5(String data)
      {
      java.security.MessageDigest digest; // Message digest of type, MD5.
      byte[] hash; // Byte array containing passed text converted to hash value.
      
      // Create a message digest with the specified algorithm name -- MD5.
      digest = java.security.MessageDigest.getInstance("MD5"); 
      
      // Convert passed text to hash value -- as a byte array.
      hash = digest.digest(data.getBytes("UTF-8"));
      
      // Return the hash value converted to a string.
      return javax.xml.bind.DatatypeConverter.printHexBinary(hash);
      }
      

      这里是另一个代码示例的链接:https://community.talend.com/t5/Design-and-Development/sha1-hash-key/td-p/109750

      比较:使用散列(可能还有其他)字段比较新旧信息。对关键字段(或相关哈希)使用完全外连接,对非关键字段使用哈希。当在新端而非旧端发现空值时,需要插入。当在旧端而不是新端发现空值时,需要删除(或根据您的需要简单地忽略)。当没有出现 null 时,执行更新。

      【讨论】:

        【解决方案4】:

        这里是一个例子:

        要求:

        • Dim_Control(Job_Id、Job_Name、Table_Name、Last_Success、Created_Date)

          CREATE TABLE Dim_Control(
              job_id BIGINT IDENTITY(1,1) PRIMARY KEY
              ,job_name NVARCHAR(255)
              ,table_name NVARCHAR(255)
              ,last_success DATETIME2(0)
              ,created_date DATETIME2(0) DEFAULT GETDATE()
          )
          
        • 上下文(Last_Success、Job_Name、Table_Name、Current_Run)

        步骤:

        1.获取上次成功的工作名称和日期:

        "Select job_name, table_name, MAX(last_success) as last_success
        FROM Dim_Control
        WHERE table_name ='Employee'
        GROUP BY job_name, table_name;"
        

        2.写日志: 使用 tLogRow 组件 - 您可以选择表格(在表格的单元格中打印值)

        3.匹配上下文变量和Dim_Control值

        System.out.println("*** Job Name = "+input_row.job_name);
        System.out.println("*** Table Name = "+input_row.table_name);
        System.out.println("*** Last Success = "+input_row.last_success);
        System.out.println("*** (Before) context last_success:" +context.last_success);
        
        context.last_success = TalendDate.formatDate("yyyy-MM-dd HH:mm:ss",input_row.last_success);
        context.current_run = TalendDate.formatDate("yyyy-MM-dd HH:mm:ss",TalendDate.getCurrentDate());
        context.table_name = input_row.table_name;
            
        System.out.println("*** (After) context last_success:" +context.last_success);
        System.out.println("*** (After) context current_run:" +context.current_run);
        

        4.截断目标阶段表

        5.向目标阶段表插入新记录:

            "SELECT distinct *
            FROM dbo.Source_Employee a WITH(NOLOCK)
            WHERE FORMAT(ISNULL(a.UpdateDate, a.CreatedDate),'yyyy-MM-dd HH:mm:ss') >=  '" + context.last_success +"' OPTION (MAXDOP 32);"
        

        6.向 Dim_Control 插入新的成功作业信息

        "INSERT INTO Dim_Control (job_name, table_name, last_success)
        VALUES ('"+context.job_name+"',  '"+context.table_name+"', '"+context.current_run+"' ); "
        

        7.合并阶段和主目标表

        "MERGE dbo.Main_Target_Table t1 
        USING dbo.Stage_Target_Table t2
        ON t1.Id = t2.Id
        WHEN MATCHED
            THEN UPDATE SET Id = t2.Id, Name= t2.Name
        WHEN NOT MATCHED BY TARGET
        THEN INSERT ( Id, Name ) VALUES ( t2.Id, t2.Name);"
        
        • 合并不应包括删除部分,ID 应为主键
        • 所有上下文类型都是字符串
        • 默认 Talend 日期格式:2021 年 6 月 28 日星期四 00:00:00 EET
        • 也可以看Rohan's Video

        Workflow of ETL

        【讨论】:

          猜你喜欢
          • 2019-06-26
          • 2015-04-22
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2021-03-31
          • 1970-01-01
          • 2022-12-14
          • 1970-01-01
          相关资源
          最近更新 更多