コンテンツにスキップ

Pythond コレクターのベストプラクティス クイックスタート


Pythond は、ユーザーが定義した Python 収集スクリプトを定期的にトリガーするための完全なソリューションです。本記事では、「1 時間ごとのログインユーザー数」を取得してメトリクスとしてセンターに報告するケースを例に説明します。

業務デモの紹介

業務の流れはおおよそ次のとおりです。
データベースからデータを収集(Python スクリプト)→ Pythond コレクターが定期的にスクリプトをトリガーしてデータを報告(Datakit)→ センターでメトリクスを確認(Web)

データベースには customers というテーブルがあり、以下のフィールドが含まれています。

  • name: 名前(文字列)
  • last_logined_time : ログイン時間(タイムスタンプ)

テーブル作成 SQL は次のとおりです。

create table customers
(
  `id`                BIGINT(20)  not null AUTO_INCREMENT COMMENT '自增 ID',
  `last_logined_time` BIGINT(20)  not null DEFAULT 0      COMMENT '登录时间 (时间戳)',
  `name`              VARCHAR(48) not null DEFAULT ''     COMMENT '姓名',

  primary key(`id`),
  key idx_last_logined_time(last_logined_time)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

上記のテーブルにテストデータを挿入します。

INSERT INTO customers (id, last_logined_time, name) VALUES (1, 1645600127, 'zhangsan');
INSERT INTO customers (id, last_logined_time, name) VALUES (2, 1645600127, 'lisi');
INSERT INTO customers (id, last_logined_time, name) VALUES (3, 1645600127, 'wangwu');

次の SQL を使用して「1 時間ごとのログインユーザー数」を取得します。

select count(1) from customers where last_logined_time>=(unix_timestamp()-3600);

上記のデータをメトリクスとしてセンターに報告します。

以下では、上記の業務を実現する具体的な手順を詳しく説明します。

前提条件

Python 環境

現在はアルファ段階であり、Python 3+ のみ互換性があります。テスト済みのバージョンは次のとおりです。

  • 3.10.1

Python 依存ライブラリ

以下の依存ライブラリをインストールする必要があります。

  • requests(ネットワーク操作、メトリクス報告に使用)
  • pymysql(MySQL データベース操作、データベース接続による業務データ取得に使用)

インストール方法は次のとおりです。

# python3
python3 -m pip install requests
python3 -m pip install pymysql

上記のインストールには pip が必要です。pip がない場合は、以下の方法を参考にしてください(出典: こちら)。

# Linux/MacOS
python3 -m ensurepip --upgrade

# Windows
py -m ensurepip --upgrade

インストールとデプロイ

1 ユーザー定義スクリプトの作成

ユーザーは DataKitFramework クラスを継承し、run メソッドをオーバーライドする必要があります。DataKitFramework クラスのソースコードファイルは datakit_framework.py で、パスは datakit/python.d/core/datakit_framework.py です。

具体的な使い方は、ソースコードファイル datakit/python.d/core/demo.py を参照してください。

上記の要件に基づいて、次のような Python スクリプトを作成し、hellopythond.py と名付けます。

hellopythond.py
from datakit_framework import DataKitFramework
import pymysql
import re
import logging

class MysqlConn():
    def __init__(self, logger, config):
        self.logger = logger
        self.config = config
        self.re_errno = re.compile(r'^\((\d+),')

        try:
            self.conn = pymysql.Connect(**self.config)
            self.logger.info("pymysql.Connect() ok, {0}".format(id(self.conn)))
        except Exception as e:
            raise e

    def __del__(self):
        self.close()

    def close(self):
        if self.conn:
            self.logger.info("conn.close() {0}".format(id(self.conn)))
            self.conn.close()


    def execute_query(self, sql_str, sql_params=(), first=True):
        res_list = None
        cur = None
        try:
            cur = self.conn.cursor()
            cur.execute(sql_str, sql_params)
            res_list = cur.fetchall()
        except Exception as e:
            err = str(e)
            self.logger.error('execute_query: {0}'.format(err))
            if first:
                retry = self._deal_with_network_exception(err)
                if retry:
                    return self.execute_query(sql_str, sql_params, False)
        finally:
            if cur is not None:
                cur.close()
        return res_list

    def execute_write(self, sql_str, sql_params=(), first=True):
        cur = None
        n = None
        err = None
        try:
            cur = self.conn.cursor()
            n = cur.execute(sql_str, sql_params)
        except Exception as e:
            err = str(e)
            self.logger.error('execute_query: {0}'.format(err))
            if first:
                retry = self._deal_with_network_exception(err)
                if retry:
                    return self.execute_write(sql_str, sql_params, False)
        finally:
            if cur is not None:
                cur.close()
        return n, err

    def _deal_with_network_exception(self, stre):
        errno_str = self._get_errorno_str(stre)
        if errno_str != '2006' and errno_str != '2013' and errno_str != '0':
            return False
        try:
            self.conn.ping()
        except Exception as e:
            return False
        return True

    def _get_errorno_str(self, stre):
        searchObj = self.re_errno.search(stre)
        if searchObj:
            errno_str = searchObj.group(1)
        else:
            errno_str = '-1'
        return errno_str

    def _is_duplicated(self, stre):
        errno_str = self._get_errorno_str(stre)
        # 1062:字段值重复,入库失败
        # 1169:字段值重复,更新记录失败
        if errno_str == "1062" or errno_str == "1169":
            return True
        return False

class HelloPythond(DataKitFramework):
    __name = 'HelloPythond'
    interval = 10 # 每 10 秒钟采集上报一次。这个根据实际业务进行调节,这里仅作演示。

    # if your datakit ip is 127.0.0.1 and port is 9529, you won't need use this,
    # just comment it.
    # def __init__(self, **kwargs):
    #     super().__init__(ip = '127.0.0.1', port = 9529)

    def run(self):
        config = {
            "host": "172.16.2.203",
            "port": 30080,
            "user": "root",
            "password": "Kx2ADer7",
            "db": "df_core",
            "autocommit": True,
            # "cursorclass": pymysql.cursors.DictCursor,
            "charset": "utf8mb4"
        }

        mysql_conn = MysqlConn(logging.getLogger(''), config)
        query_str = "select count(1) from customers where last_logined_time>=(unix_timestamp()-%s)"
        sql_params = ('3600')
        n = mysql_conn.execute_query(query_str, sql_params)

        data = [
        {
            "measurement": "hour_logined_customers_count", # 指标名称。
            "tags": {
                "tag_name": "tag_value", # 自定义 tag,根据自己想要标记的填写,我这里是随便写的
            },
            "fields": {
                "count": n[0][0], # 指标,这里是每个小时登录的用户数
            },
        },
        ]

        in_data = {
            'M':data,
            'input': "pyfromgit"
        }

        return self.report(in_data) # you must call self.report here

2 カスタムスクリプトを適切な場所に配置する

Datakit のインストールディレクトリにある python.d ディレクトリの下に新しいフォルダを作成し、hellopythond という名前にします。このフォルダ名は、上記で作成したクラス名と同じ(hellopythond)にする必要があります。

次に、作成したスクリプト hellopythond.py をこのフォルダに配置します。最終的なディレクトリ構造は次のとおりです。

├── ...
├── datakit
└── python.d
    ├── core
    │   ├── datakit_framework.py
    │   └── demo.py
    └── hellopythond
        └── hellopythond.py

注意: 上記の core フォルダは Pythond のコアフォルダです。変更しないでください。

上記は gitrepos 機能を有効にしていない場合の構成です。gitrepos 機能を有効にしている場合のパス構造は次のとおりです。

├── ...
├── datakit
├── python.d
├── gitrepos
│   └── yourproject
│       ├── conf.d
│       ├── pipeline
│       └── python.d
│           └── hellopythond
│               └── hellopythond.py

3 Pythond 設定ファイルを有効にする

Pythond 設定ファイルをコピーします。
conf.d/pythond ディレクトリで pythond.conf.samplepythond.conf としてコピーし、次のように設定します。

[[inputs.pythond]]

    # Python コレクター名
    name = 'some-python-inputs'  # required

    # Python コレクターの実行に必要な環境変数
    #envs = ['LD_LIBRARY_PATH=/path/to/lib:$LD_LIBRARY_PATH',]

    # Python コレクターの実行可能プログラムのパス(可能な限り絶対パスを記述してください)
    cmd = "python3" # required. python3 is recommended.

    # ユーザースクリプトの相対パス(フォルダ名を指定します。指定したフォルダの直下のモジュールと py ファイルがすべて適用されます)
    dirs = ["hellopythond"] # ここにはフォルダ名(クラス名)を入力します

4 DataKit を再起動する

sudo datakit --restart

動作確認

すべてが正常に動作すれば、約 1 分以内にセンターでメトリクスグラフを確認できます。

14.pythond.png

参考資料

<公式マニュアル: Python を使用したカスタムコレクターの開発>

<公式マニュアル: Git による設定ファイルの管理>

フィードバック

このページは役に立ちましたか?