从关系型数据库摄取数据
本指南说明了如何使用 Logstash 和 Logstash JDBC 输入插件将数据从关系型数据库提取到 Elastic Cloud 中。它演示了如何使用 Logstash 高效地复制记录并从关系型数据库接收更新,然后将其发送到 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 部署中的 Elasticsearch。
此处介绍的代码和方法已在 MySQL 上进行了测试。它们也应该适用于其他关系型数据库。
Logstash Java 数据库连接 (JDBC) 输入插件使您能够从 MySQL 和 Postgres 等许多流行的关系型数据库中提取数据。从概念上讲,JDBC 输入插件运行一个循环,定期轮询关系型数据库,查找自上次循环迭代以来插入或修改的记录。
所需时间:2 小时
对于本教程,您需要一个供 Logstash 读取的源 MySQL 实例。您可以从 MySQL 社区下载站点的“MySQL 社区服务器”部分获取免费版本的 MySQL。
- 获取免费试用.
- 登录 Elastic Cloud。
- 选择创建部署 (Create deployment)。
- 为您的部署命名。您可以保留所有其他设置的默认值。
- 选择 Create deployment(创建部署)并保存您的 Elastic 部署凭据。稍后您将需要这些凭据。
- 当部署就绪后,点击 Continue(继续),系统将显示 Setup guides(设置指南)页面。要继续前往部署主页,请点击 I’d like to do something else(我想做其他事情)。
不想订阅其他服务?您也可以通过 AWS、Azure 和 GCP 市场获取 Elastic Cloud Hosted。
- 登录 Elastic Cloud Enterprise 管理控制台。
- 选择创建部署 (Create deployment)。
- 为您的部署命名。您可以保留所有其他设置的默认值。
- 选择 Create deployment(创建部署)并保存您的 Elastic 部署凭据。稍后您将需要这些凭据。
- 当部署就绪后,点击 Continue(继续),系统将显示 Setup guides(设置指南)页面。要继续前往部署主页,请点击 I’d like to do something else(我想做其他事情)。
连接到 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 时,您可以使用云 ID (Cloud ID) 来指定连接详细信息。前往 Kibana 主菜单并选择“管理 (Management) > 集成 (Integrations)”,然后选择“查看部署详情 (View deployment details)”,即可找到您的云 ID。
为了连接、向其流式传输数据以及发出查询,您需要考虑身份验证。支持两种身份验证机制:API 密钥和基本身份验证。为了让您快速上手,我们将在此展示如何使用基本身份验证,但您也可以按照稍后的说明生成 API 密钥。API 密钥更安全,是生产环境的首选。
- 下载并在托管 MySQL 的本地机器或其他被授予访问 MySQL 机器权限的机器上解压 Logstash。
Logstash JDBC 输入插件不包含任何数据库连接驱动程序。您需要为您的关系型数据库准备一个 JDBC 驱动程序,用于后续配置带有 JDBC 输入插件的 Logstash 管道章节中的步骤。
- 从 MySQL 社区下载站点的“Connector/J”部分下载并解压 MySQL 的 JDBC 驱动程序。
- 记下驱动程序的位置,后续步骤中会用到它。
让我们来看一个简单的数据库,您将从中导入数据并将其发送到 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 部署。此示例使用带有带时间戳记录的 MySQL 数据库。时间戳使您能够轻松确定自上次数据传输以来数据库中发生了什么变化。
在此示例中,让我们创建一个新数据库 es_db 和表 es_table,作为我们 Elasticsearch 数据的来源。
运行以下 SQL 语句以生成一个包含三列表的新 MySQL 数据库
CREATE DATABASE es_db; USE es_db; DROP TABLE IF EXISTS es_table; CREATE TABLE es_table ( id BIGINT(20) UNSIGNED NOT NULL, PRIMARY KEY (id), UNIQUE KEY unique_id (id), client_name VARCHAR(32) NOT NULL, modification_time TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP );让我们探讨一下此 SQL 代码片段中的关键概念
- es_table
- 存储数据的表名。
- id
- 记录的唯一标识符。id 被定义为 PRIMARY KEY 和 UNIQUE KEY,以确保每个 id 在当前表中仅出现一次。在将文档更新或插入 Elasticsearch 时,这会转换为 _id。
- client_name
- 最终将被摄入 Elasticsearch 的数据。为简单起见,此示例仅包含一个数据字段。
- modification_time
-
记录插入或最后更新的时间戳。稍后,您可以使用此时间戳来确定自上次向 Elasticsearch 传输数据以来发生了哪些变化。
考虑如何处理删除操作以及如何通知 Elasticsearch。通常,删除记录会导致其立即从 MySQL 数据库中移除。没有关于该删除的记录。Logstash 无法检测到此更改,因此该记录仍保留在 Elasticsearch 中。
有两种可能的解决方法
- 您可以在源数据库中使用“软删除”。本质上,记录首先通过布尔标志标记为已删除。当前使用您源数据库的其他程序必须在查询中过滤掉“软删除”。“软删除”会被发送到 Elasticsearch,并在那里进行处理。之后,您的源数据库和 Elasticsearch 都必须移除这些“软删除”。
- 您可以定期清除基于数据库的 Elasticsearch 索引,然后用数据库内容的全新摄入来刷新 Elasticsearch。
登录到您的 MySQL 服务器并向新数据库添加三条记录
use es_db INSERT INTO es_table (id, client_name) VALUES (1,"Targaryen"), (2,"Lannister"), (3,"Stark");使用 SQL 语句验证您的数据
select * from es_table;输出应该类似于以下内容
+----+-------------+---------------------+ | id | client_name | modification_time | +----+-------------+---------------------+ | 1 | Targaryen | 2021-04-21 12:17:16 | | 2 | Lannister | 2021-04-21 12:17:16 | | 3 | Stark | 2021-04-21 12:17:16 | +----+-------------+---------------------+现在,让我们回到 Logstash 并将其配置为摄入此数据。
让我们设置一个示例 Logstash 输入管道,从您新的 JDBC 插件和 MySQL 数据库中摄入数据。除 MySQL 外,您还可以从任何支持 JDBC 的数据库中输入数据。
在
<localpath>/logstash-7.12.0/中,创建一个名为jdbc.conf的新文本文件。复制并粘贴以下代码到这个新文本文件中。此代码通过 JDBC 插件创建了一个 Logstash 管道。
input { jdbc { jdbc_driver_library => "<driverpath>/mysql-connector-java-<versionNumber>.jar" jdbc_driver_class => "com.mysql.jdbc.Driver" jdbc_connection_string => "jdbc:mysql://<MySQL host>:3306/es_db" jdbc_user => "<myusername>" jdbc_password => "<mypassword>" jdbc_paging_enabled => true tracking_column => "unix_ts_in_secs" use_column_value => true tracking_column_type => "numeric" schedule => "*/5 * * * * *" statement => "SELECT *, UNIX_TIMESTAMP(modification_time) AS unix_ts_in_secs FROM es_table WHERE (UNIX_TIMESTAMP(modification_time) > :sql_last_value AND modification_time < NOW()) ORDER BY modification_time ASC" } } filter { mutate { copy => { "id" => "[@metadata][_id]"} remove_field => ["id", "@version", "unix_ts_in_secs"] } } output { stdout { codec => "rubydebug"} }- 指定您的本地 JDBC 驱动程序 .jar 文件的完整路径(包括版本号)。例如:
jdbc_driver_library => "/usr/share/mysql/mysql-connector-java-8.0.24.jar" - 提供您 MySQL 主机的 IP 地址或主机名以及端口。例如,
jdbc_connection_string => "jdbc:mysql://127.0.0.1:3306/es_db" - 提供您的 MySQL 凭据。用户名和密码都必须用引号括起来。
注意如果您使用的是 MariaDB(MySQL 的一个流行的开源社区分支),则有几件事需要以不同方式处理
使用 MariaDB 的 JDBC 驱动程序替换 MySQL JDBC 驱动程序。
在
jdbc.conf代码中替换以下行,包括最后一行中的ANSI_QUOTES片段jdbc_driver_library => "<driverPath>/mariadb-java-client-<versionNumber>.jar" jdbc_driver_class => "org.mariadb.jdbc.Driver" jdbc_connection_string => "jdbc:mariadb://<mySQLHost>:3306/es_db?sessionVariables=sql_mode=ANSI_QUOTES"以下是关于 Logstash 管道代码的一些额外详细信息
- jdbc_driver_library
- Logstash JDBC 插件不随附 JDBC 驱动程序库。必须使用
jdbc_driver_library配置选项将 JDBC 驱动程序库显式传递给插件。 - tracking_column
- 此参数指定了字段
unix_ts_in_secs,用于跟踪 Logstash 从 MySQL 读取的最后一份文档,并存储在磁盘上的 logstash_jdbc_last_run 中。该参数决定了 Logstash 在轮询循环的下一次迭代中请求文档的起始值。存储在logstash_jdbc_last_run中的值可以在 SELECT 语句中作为sql_last_value访问。 - unix_ts_in_secs
- 由 SELECT 语句生成的字段,其中包含作为标准 Unix 时间戳(自 epoch 以来的秒数)的
modification_time。该字段由tracking column引用。使用 Unix 时间戳而不是普通时间戳来跟踪进度,因为普通时间戳可能会因在 UMT 和本地时区之间正确转换的复杂性而导致错误。 - sql_last_value
- 这是一个 内置参数,包含 Logstash 轮询循环当前迭代的起点,并在 JDBC 输入配置的 SELECT 语句行中被引用。此参数设置为从
.logstash_jdbc_last_run读取的unix_ts_in_secs的最新值。此值是 Logstash 轮询循环中执行的 MySQL 查询所返回文档的起点。在查询中包含此变量可确保我们不会重新发送已存储在 Elasticsearch 中的数据。 - schedule
- 这使用 cron 语法来指定 Logstash 轮询 MySQL 更改的频率。规范
*/5 * * * * *告诉 Logstash 每 5 秒联系一次 MySQL。此插件的输入可以安排根据特定计划定期运行。这种调度语法由 rufus-scheduler 提供支持。该语法类似于 cron,并带有特定于 Rufus 的一些扩展(例如,时区支持)。 - modification_time < NOW()
- SELECT 的这部分内容将在下一节中详细解释。
- filter
- 在本节中,
id的值从 MySQL 记录复制到名为_id的元数据字段中,该字段稍后在输出中被引用,以确保每个文档都以正确的_id值写入 Elasticsearch。使用元数据字段可确保此临时值不会导致创建新字段。id、@version和unix_ts_in_secs字段也会从文档中删除,因为它们不需要写入 Elasticsearch。 - 输出 (output)
-
本节指定每个文档都应使用 rubydebug 输出写入标准输出,以帮助进行调试。
- 指定您的本地 JDBC 驱动程序 .jar 文件的完整路径(包括版本号)。例如:
使用新的 JDBC 配置文件启动 Logstash
bin/logstash -f jdbc.confLogstash 通过标准输出 (
stdout)(即您的命令行界面)输出您的 MySQL 数据。初始数据加载的结果应该类似于以下内容[INFO ] 2021-04-21 12:32:32.816 [Ruby-0-Thread-15: :1] jdbc - (0.009082s) SELECT * FROM (SELECT *, UNIX_TIMESTAMP(modification_time) AS unix_ts_in_secs FROM es_table WHERE (UNIX_TIMESTAMP(modification_time) > 0 AND modification_time < NOW()) ORDER BY modification_time ASC) AS 't1' LIMIT 100000 OFFSET 0 { "client_name" => "Targaryen", "modification_time" => 2021-04-21T12:17:16.000Z, "@timestamp" => 2021-04-21T12:17:16.923Z } { "client_name" => "Lannister", "modification_time" => 2021-04-21T12:17:16.000Z, "@timestamp" => 2021-04-21T12:17:16.961Z } { "client_name" => "Stark", "modification_time" => 2021-04-21T12:17:16.000Z, "@timestamp" => 2021-04-21T12:17:16.963Z }即使 MySQL 数据库中没有新的或修改的内容,Logstash 结果也会定期显示 SQL SELECT 语句
[INFO ] 2021-04-21 12:33:30.407 [Ruby-0-Thread-15: :1] jdbc - (0.002835s) SELECT count(*) AS 'count' FROM (SELECT *, UNIX_TIMESTAMP(modification_time) AS unix_ts_in_secs FROM es_table WHERE (UNIX_TIMESTAMP(modification_time) > 1618935436 AND modification_time < NOW()) ORDER BY modification_time ASC) AS 't1' LIMIT 1打开您的 MySQL 控制台。让我们使用以下 SQL 语句向该数据库中插入另一条记录
use es_db INSERT INTO es_table (id, client_name) VALUES (4,"Baratheon");切回您的 Logstash 控制台。Logstash 会检测到新记录,控制台会显示类似于以下内容的结果
[INFO ] 2021-04-21 12:37:05.303 [Ruby-0-Thread-15: :1] jdbc - (0.001205s) SELECT * FROM (SELECT *, UNIX_TIMESTAMP(modification_time) AS unix_ts_in_secs FROM es_table WHERE (UNIX_TIMESTAMP(modification_time) > 1618935436 AND modification_time < NOW()) ORDER BY modification_time ASC) AS 't1' LIMIT 100000 OFFSET 0 { "client_name" => "Baratheon", "modification_time" => 2021-04-21T12:37:01.000Z, "@timestamp" => 2021-04-21T12:37:05.312Z }检查 Logstash 输出结果,确保您的数据看起来正确。使用
CTRL + C关闭 Logstash。
在本节中,我们将配置 Logstash 将 MySQL 数据发送到 Elasticsearch。我们修改在 配置带有 JDBC 输入插件的 Logstash 管道 章节中创建的配置文件,以便将数据直接输出到 Elasticsearch。我们启动 Logstash 发送数据,然后登录到您的部署以验证 Kibana 中的数据。
打开 Logstash 文件夹中的
jdbc.conf文件进行编辑。用以下内容更新输出部分
output { elasticsearch { index => "rdbms_idx" ilm_enabled => false cloud_id => "<DeploymentName>:<ID>" cloud_auth => "elastic:<Password>" # api_key => "<myAPIid:myAPIkey>" } }- 使用您的 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 部署的云 ID。您可以在云 ID 开头包含或省略
<DeploymentName>:前缀。两个版本都可以正常工作。前往 Kibana 主菜单并选择“管理 (Management) > 集成 (Integrations)”,然后选择“查看部署详情 (View deployment details)”,即可找到您的云 ID。 - 默认用户名为
elastic。不建议使用elastic账户来摄入数据,因为这是超级用户。我们建议使用权限受限的用户,或具有特定于要写入的索引或数据流权限的 API 密钥。有关角色和 API 密钥的信息,请查看 在 Logstash 中配置安全性。如果使用elastic用户,请使用创建部署时提供的密码;如果使用 在 Logstash 中配置安全性 文档中指定了角色的新摄入用户,则使用创建该用户时使用的密码。
以下是关于配置文件设置的一些额外详细信息
- index
- Elasticsearch 索引名称
rdbms_idx,用于关联文档。 - api_key
-
如果您选择使用 API 密钥进行身份验证(在下一步中讨论),您可以在此处提供。
- 使用您的 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 部署的云 ID。您可以在云 ID 开头包含或省略
可选:为了增加安全性,您可以通过 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 控制台生成一个 Elasticsearch API 密钥,并配置 Logstash 使用该新密钥安全连接到您的部署。
对于 Elastic Cloud Hosted,登录 Elastic Cloud;对于 Elastic Cloud Enterprise,登录管理控制台。
选择部署名称并前往 ☰ > 管理 (Management) > 开发工具 (Dev Tools)。
输入以下内容
POST /_security/api_key { "name": "logstash-apikey", "role_descriptors": { "logstash_read_write": { "cluster": ["manage_index_templates", "monitor"], "index": [ { "names": ["logstash-*","rdbms_idx"], "privileges": ["create_index", "write", "read", "manage"] } ] } } }这将创建一个具有集群
monitor权限的 API 密钥,该权限提供用于确定集群状态的只读访问权限,并且manage_index_templates允许对索引模板执行所有操作。一些额外的权限还允许对指定索引执行create_index、write和manage操作。添加索引manage权限是为了启用索引刷新。点击 ▶。输出应该类似于以下内容
{ "api_key": "tV1dnfF-GHI59ykgv4N0U3", "id": "2TBR42gBabmINotmvZjv", "name": "logstash_api_key" }将您的新
api_key值输入到 Logstashjdbc.conf文件中,格式为<id>:<api_key>。如果您的结果如本示例所示,您将输入2TBR42gBabmINotmvZjv:tV1dnfF-GHI59ykgv4N0U3。记得删除井号 (#) 以取消注释该行,并注释掉username和password行output { elasticsearch { index => "rdbms_idx" cloud_id => "<myDeployment>" ilm_enabled => false api_key => "2TBR42gBabmINotmvZjv:tV1dnfF-GHI59ykgv4N0U3" # user => "<Username>" # password => "<Password>" } }
如果您只是按照原样重启 Logstash 并使用新的输出,那么将不会有 MySQL 数据发送到我们的 Elasticsearch 索引。
为什么?Logstash 保留了之前的
sql_last_value时间戳,并发现自那时以来 MySQL 数据库中没有发生新的更改。因此,基于我们配置的 SQL 查询,没有新数据要发送到 Logstash。解决方案:在
jdbc.conf文件的 JDBC 输入部分添加clean_run => true作为新行。当设置为true时,此参数会将sql_last_value重置为零。input { jdbc { ... clean_run => true ... } }在将
sql_last_value设置为true运行 Logstash 一次后,您可以删除clean_run行,除非您希望在每次重启 Logstash 时都发生重置行为打开命令行界面实例,转到您的 Logstash 安装路径,然后启动 Logstash
bin/logstash -f jdbc.confLogstash 将 MySQL 数据输出到您的 Elastic Cloud Hosted 或 Elastic Cloud Enterprise 部署。让我们在 Kibana 中查看并验证该数据
对于 Elastic Cloud Hosted,登录 Elastic Cloud;对于 Elastic Cloud Enterprise,登录管理控制台。
选择部署并前往 ☰ > 管理 (Management) > 开发工具 (Dev Tools)
将以下 API GET 请求复制并粘贴到控制台窗格中,然后点击 ▶。这将查询新
rdbms_idx索引中的所有记录。GET rdbms_idx/_search { "query": { "match_all": {} } }结果窗格会列出源自您的 MySQL 数据库的
client_name记录,类似于以下示例
现在,您应该对如何通过 JDBC 插件配置 Logstash 以从关系型数据库中摄入数据有了很好的了解。您还需要考虑一些设计因素来跟踪新增、修改和删除的记录。您现在应该已经掌握了开始试验您自己的数据库和 Elasticsearch 所需的基础知识。