通过DataWorks将MaxCompute数据同步到阿里云ES

更新时间:
复制 MD 格式

如果您需要对MaxCompute(ODPS)中的海量数据进行信息检索、多维查询、统计分析等操作,可借助阿里云Elasticsearch实现。本文通过DataWorks的数据集成服务,实现最快分钟级,将海量MaxCompute数据离线同步到阿里云ES中。

背景信息

DataWorks是一个基于大数据引擎,集成数据开发、任务调度、数据管理等功能的全链路大数据开发治理平台。您可以通过DataWorks的同步任务,快速的将各种数据源中的数据同步到阿里云ES。

  • 支持同步的数据源包括:

    • 阿里云云数据库(MySQL、PostgreSQL、SQL Server、MongoDB、HBase)

    • 阿里云PolarDB-X(原DRDS升级版)

    • 阿里云MaxCompute

    • 阿里云OSS

    • 阿里云Tablestore

    • 自建HDFS、Oracle、FTP、DB2及以上数据库类型的自建版本

  • 适用场景:

前提条件

说明
  • 仅支持将数据同步到阿里云ES,不支持自建Elasticsearch。

  • MaxCompute项目、ES实例和DataWorks工作空间所在地域需保持一致。

  • ES实例、MaxComputeDataWorks工作空间需要创建在同一时区下,否则同步与时间相关的数据时,同步前后的数据可能存在时区差。

费用说明

操作步骤

步骤一:准备源数据

创建MaxCompute表并导入测试数据。具体操作,请参见创建表导入数据

本文使用的表结构和表数据如下所示:

  • 表结构

    表包含 7 个字段和 1 个分区字段,字段定义如下:

    • create_time(string,主键)

    • category(string)

    • brand(string)

    • buyer_id(string)

    • trans_num(bigint)

    • trans_amount(double)

    • click_cnt(bigint)

    分区字段为 pt(bigint)。

  • 部分表数据

    源数据表包含以下字段:

    • create_time:交易日期,如 2020/6/1

    • category:商品类目,如外套、生鲜、电器、卫浴

    • brand:品牌名称,如品牌 A~品牌 G

    • buyer_id:买家 ID,如 user1~user13

    • trans_num:交易数量

    • trans_amount:交易金额

    • click_cnt:点击次数

    • pt:分区字段,值为 1

步骤二:购买并配置独享资源组

购买一个数据集成独享资源组,并为该资源组绑定专有网络和工作空间。独享资源组可以保证数据快速、稳定地传输。

  1. 登录DataWorks控制台

  2. 在顶部菜单栏选择相应地域后,在左侧导航栏单击资源组

  3. 独享资源组页签下,单击创建旧版资源组 > 创建旧版资源组 > 数据集成资源组

  4. DataWorks独享资源(包年包月)购买页面,独享资源类型选择独享数据集成资源,输入资源组名称,单击立即购买,购买独享资源组。

    更多配置信息,请参见步骤一:购买资源组

  5. 在已创建的独享资源组的操作列,单击网络设置,为该独享资源组绑定专有网络。具体操作,请参见绑定专有网络

    说明

    本文以独享数据集成资源组通过VPC内网同步数据为例。关于通过公网同步数据,请参见添加白名单

    独享资源需要与Elasticsearch实例的专有网络连通才能同步数据。因此需要绑定Elasticsearch实例所在的专有网络可用区交换机。查看Elasticsearch实例所在的专有网络、可用区和交换机,请参见查看Elasticsearch实例的基本信息

    重要

    绑定专有网络后,您需要将专有网络的交换机网段加入到Elasticsearch实例的VPC私网访问白名单中。具体操作,请参见配置Elastic search实例公网或私网访问白名单

  6. 在页面左上角,单击返回图标,返回资源组列表页面。

  7. 在已创建的独享资源组的操作列,单击绑定工作空间,为该独享资源组绑定目标工作空间。

    具体操作,请参见步骤二:绑定归属工作空间

步骤三:添加数据源

MaxComputeElasticsearch数据源接入DataWorks的数据集成服务中。

  1. 进入DataWorks快速进入 > 数据集成页面。

    1. 登录DataWorks控制台

    2. 在左侧导航栏,单击工作空间

    3. 在目标工作空间的操作列,选择快速进入 > 快速进入 > 数据集成

  2. 在左侧导航栏,单击数据源

  3. 新增MaxCompute数据源。

    1. 数据源列表页面,单击新增数据源

    2. 新增数据源页面,搜索并选择MaxCompute数据源。

    3. 新增MaxCompute数据源对话框,在基础信息区域配置数据源参数。

      配置详情,请参见配置MaxCompute数据源

    4. 连接配置区域,单击测试连通性,连通状态显示为可连通时,表示连通成功。

    5. 单击完成

  4. 使用同样的方式添加Elasticsearch数据源。配置详情,请参见配置Elasticsearch数据源

步骤四:配置并运行数据同步任务

数据同步任务将独享资源组作为一个可以执行任务的资源,独享资源组将获取数据集成服务中数据源的数据,并将数据写入Elasticsearch。

说明
  1. 进入DataWorks快速进入 > 数据集成。 在左侧导航栏,单击数据源页面。

    1. 登录DataWorks控制台

    2. 在左侧导航栏,单击工作空间

    3. 在目标工作空间的操作列,选择快速进入 > 快速进入 > 数据集成。 在左侧导航栏,单击数据源

  2. 新建一个离线同步任务。

    1. 在左侧导航栏的数据开发(image图标)页签,选择 新建 > 新建业务流,按照界面指引新建一个业务流程。

    2. 右键单击新建的业务流程,选择新建节点 > 新建节点 > 数据集成 > 离线同步 > 到阿里云ES

    3. 新建节点对话框中,输入节点名称,单击确认

  3. 配置网络与资源

    1. 数据来源区域,数据来源选择MaxCompute(ODPS),数据源名称选择待同步的数据源名称。 在我的资源组区域,数据来源区域,数据来源选择MaxCompute(ODPS),数据源名称选择待同步的数据源名称。 在我的资源组选择MaxCompute(ODPS),数据源名称选择待同步的数据源名称。

    2. 我的资源组区域,选择独享资源组。

    3. 数据去向区域,数据去向选择Elasticsearch,数据源名称选择待同步的数据源名称。

  4. 单击下一步

  5. 配置任务。

    1. 数据来源区域,数据来源选择MaxCompute(ODPS),数据源名称选择待同步的数据源名称。 在我的资源组区域,选择待同步的表。

    2. 数据去向区域,配置数据去向的各参数。

    3. 字段映射区域中,设置来源字段目标字段的映射关系。

    4. 通道控制区域,配置通道参数。

    详细配置信息,请参见向导模式配置

  6. 运行任务。

    1. (可选)配置任务调度属性。在页面右侧,单击调度配置,按照需求配置相应的调度参数。各配置的详细说明,请参见调度配置

    2. 在节点区域的右上角,单击保存图标,保存任务。

    3. 在节点区域的右上角,单击提交图标,提交任务。

      如果您配置了任务调度属性,任务会定期自动执行。您还可以在节点区域的右上角,单击运行图标,立即运行任务。

      运行日志中出现Shell run successfully!表明任务运行成功。部分任务运行日志如下所示:

      2023-10-31 16:52:35 INFO Exit code of the Shell command 0
      2023-10-31 16:52:35 INFO --- Invocation of Shell command completed ---
      2023-10-31 16:52:35 INFO Shell run successfully!
      2023-10-31 16:52:35 INFO Current task status: FINISH
      2023-10-31 16:52:35 INFO Cost time is: 33.106s

步骤五:验证数据同步结果

Kibana控制台中,查看同步成功的数据,并按条件查询数据。

  1. 登录目标阿里云Elasticsearch实例的Kibana控制台。

    具体操作,请参见登录Kibana控制台

  2. 单击Kibana页面左上角的菜单.png图标,选择Dev Tools(开发工具)。

  3. Console(控制台)中,执行如下命令查看同步的数据。

    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {}}
    }
    说明

    odps_index为您在数据同步脚本中设置的index字段的值。

    数据同步成功后,返回如下结果。

    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {}}
    }
    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {} },
    "_source": ["category", "brand"]
    }
    POST /odps_index/_search?pretty
    {
      "query": { "match": {"category":"生鲜"} }
    }
    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {} },
    "sort": { "trans_num": { "order": "desc" }
            }
    }
    
    --- Response ---
    {
      "took" : 2,
      "timed_out" : false,
      "_shards" : {
        "total" : 1,
        "successful" : 1,
        "skipped" : 0,
        "failed" : 0
      },
      "hits" : {
        "total" : 13,
        "max_score" : null,
        "hits" : [
          {
            "_index" : "odps_index",
            "_type" : "_doc",
            "_id" : "2020/6/7 8:00",
            "_score" : null,
            "_source" : {
              "trans_num" : 88,
              "click_cnt" : 80,
              "category" : "外套",
              "buyer_id" : "user7",
              "trans_amount" : 150.0,
              "brand" : "品牌E"
            },
            "sort" : [
              88
            ]
          },
          {
            "_index" : "odps_index",
            "_type" : "_doc",
            "_id" : "2020/6/11 8:00",
            "_score" : null,
            "_source" : {
              "trans_num" : 22,
              "click_cnt" : 70,
              "category" : "卫浴",
              "buyer_id" : "user11",
              "trans_amount" : 4500.0,
              "brand" : "品牌G"
            },
            "sort" : [
              22
            ]
          }
        ]
      }
    }
  4. 执行如下命令,搜索文档中的categorybrand字段。

    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {} },
    "_source": ["category", "brand"]
    }
  5. 执行如下命令,搜索category生鲜的文档。

    POST /odps_index/_search?pretty
    {
    "query": { "match": {"category":"生鲜"} }
    }
  6. 执行如下命令,按照trans_num字段对文档进行排序。

    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {} },
    "sort": { "trans_num": { "order": "desc" } }
    }

    更多命令和访问方式,请参见Elastic.co官方帮助中心