すべてのプロダクト
Search
ドキュメントセンター

DataHub:Flume プラグイン

最終更新日:Mar 13, 2026

Flume-DataHub プラグインは、Flume 上に構築された DataHub 用の変更追跡およびパブリッシングプラグインです。これを使用すると、DataHub にデータを書き込んだり、DataHub からデータを読み取って他のシステムに書き込んだりできます。このプラグインは Flume 開発標準に準拠しており、インストールが簡単で、DataHub でデータをパブリッシュおよびサブスクライブできます。

Flume プラグインのインストール

インストール制限

  • JDK 1.8 以降

  • Apache Maven バージョン 3.x

  • Flume-NG バージョン 1.x

Flume のインストール

  1. Flume をダウンロードします。すでに Flume をダウンロードしている場合は、このステップをスキップできます。

    $ tar zxvf apache-flume-1.11.0-bin.tar.gz
    説明

    このドキュメントでは、${FLUME_HOME} は Flume のホームディレクトリを指します。

  2. Flume-DataHub をインストールします。

    • 直接インストール

      1. Flume-DataHub プラグインをダウンロードします。

      2. Flume プラグインを解凍し、${FLUME_HOME}/plugins.d ディレクトリに移動します。

        $ tar aliyun-flume-datahub-sink-x.x.x.tar.gz
        $ cd aliyun-flume-datahub-sink-x.x.x
        $ mkdir ${FLUME_HOME}/plugins.d
        $ mv aliyun-flume-datahub-sink ${FLUME_HOME}/plugins.d
    • ソースコードからのインストール。

      1. aliyun-maxcompute-data-collectors からソースコードをダウンロードします。

      2. コンパイルとインストール。

        $ cd aliyun-maxcompute-data-collectors
        $ mvn clean package -DskipTests=true  -Dmaven.javadoc.skip=true
        $ cd flume-plugin/target
        $ tar zxvf aliyun-flume-datahub-sink-x.x.x.tar.gz
        $ mv aliyun-flume-datahub-sink ${FLUME_HOME}/plugins.d

パラメーターリファレンス

Sink パラメーター

名前

デフォルト値

必須

説明

datahub.endPoint

-

必須

Alibaba Cloud DataHub のサービスエンドポイント。

datahub.accessId

-

必須

ご利用の Alibaba Cloud AccessKey ID。

datahub.accessKey

-

必須

Alibaba Cloud AccessKey。

datahub.project

-

必須

DataHub プロジェクトの名前。

datahub.topic

-

必須

DataHub Topic の名前。

datahub.shard.ids

すべてのシャード

オプション

DataHub に書き込む特定のシャード ID のカンマ区切りリスト (例: 0,1,2)。毎回、リストからシャードがランダムに選択され、データが書き込まれます。このパラメーターが指定されていない場合、Flume は Shard Split またはマージ後にシャードリストを自動的に調整します。それ以外の場合は、構成ファイルを手動で変更する必要があります。

datahub.enablePb

true

オプション

データ転送に Protocol Buffers (PB) を使用するかどうかを指定します。一部の Apsara Stack 環境では PB をサポートしていないため、手動で false に設定する必要があります。

datahub.compressType

none

オプション

データ転送のためにデータを圧縮するかどうかを指定します。サポートされている圧縮タイプは LZ4 および DEFLATE です。

datahub.batchSize

1000

オプション

DataHub への 1 回の送信あたりの最大データ量。

datahub.maxBufferSize

2 × 1024 × 1024

オプション

1 回のリクエストで書き込まれるデータの最大サイズ (バイト単位)。このパラメーターは変更しないでください。一度に大量のデータを書き込むと、書き込み操作が失敗する可能性があります。

datahub.batchTimeout

5

オプション

レコード数がバッチサイズに達しない場合に、DataHub にデータを同期するまでの待機時間 (秒単位)。

datahub.retryTimes

3

オプション

データ同期が失敗した場合の再試行回数。

datahub.retryInterval

5

オプション

データ同期が失敗した場合の再試行間隔 (秒単位)。

datahub.dirtyDataContinue

true

オプション

ダーティデータが検出されたときに処理を続行するかどうかを指定します。true に設定すると、ダーティデータはカンマを区切り文字としてダーティデータファイルに自動的に書き込まれます。これは、後続のデータの処理には影響しません。

datahub.dirtyDataFile

DataHub-Flume-dirty-file

オプション

ダーティデータファイル。

serializer

-

必須

データ解析方法。サポートされているメソッドは DELIMITED (区切り文字)、JSON (各行は単一レイヤーの JSON オブジェクト)、および REGEX (正規表現) です。

serializer.delimiter

,

オプション

データフィールドの区切り文字。特殊文字を使用する場合は、"\t" のように二重引用符で囲みます。

serializer.regex

(.*)

オプション

データ解析用の正規表現。各フィールドのデータはグループに解析されます。

serializer.fieldnames

-

必須

入力データフィールドから DataHub フィールドへのマッピング。フィールドは入力順序で識別されます。フィールドをスキップするには、その列名を空にします。たとえば、`c1,c2,,c3` は、入力データの最初の、2番目の、および4番目のフィールドを DataHub の `c1`、`c2`、および `c3` フィールドにマッピングします。

serializer.charset

UTF-8

オプション

データ解析のエンコード形式。

Source パラメーター

名前

デフォルト値

必須

説明

datahub.endPoint

-

必須

Alibaba Cloud DataHub のサービスエンドポイント。

datahub.accessId

-

必須

Alibaba Cloud AccessKey ID

datahub.accessKey

-

必須

Alibaba Cloud AccessKey。

datahub.project

-

必須

DataHub プロジェクトの名前。

datahub.topic

-

必須

DataHub Topic の名前。

datahub.subId

-

必須

DataHub サブスクリプション ID。

datahub.startTime

-

オプション

データ読み取りを開始する時点を指定します。フォーマットは yyyy-MM-dd HH:mm:ss です。このパラメーターを使用すると、まずサブスクリプションがリセットされ、その後サブスクリプションに基づいてデータが読み取られます。

datahub.shard.ids

-

オプション

DataHub から読み取る特定のシャード ID のカンマ区切りリスト (例: 0,1,2)。データが読み取られるたびに、リストからシャードがランダムに選択され、消費されます。このパラメーターが指定されていない場合、協調消費がデータ読み取りに使用されます。このパラメーターは使用しないでください。複数のソースが構成されており、このパラメーターが指定されていない場合、協調消費はソース間の負荷分散を確保するためにシャードを自動的に割り当てます。

datahub.enablePb

true

オプション

データ転送に PB を使用するかどうかを指定します。一部の Apsara Stack 環境では PB をサポートしていないため、手動で false に設定する必要があります。

datahub.compressType

none

オプション

データ転送のためにデータを圧縮するかどうかを指定します。サポートされている圧縮タイプは LZ4 および DEFLATE です。

datahub.batchSize

1000

オプション

一度に DataHub から読み取るレコードの最大数。

datahub.batchTimeout

5

オプション

レコード数がバッチサイズに達しない場合に、DataHub にデータを同期するまでの待機時間 (秒単位)。

datahub.retryTimes

3

オプション

データ読み取りが失敗した場合の再試行回数。再試行間隔は 1 秒に固定されており、調整できません。

datahub.autoCommit

true

オプション

true に設定すると、コンシューマーは自動的にオフセットをコミットします。これにより、データが消費される前にオフセットがコミットされる可能性があります。false に設定すると、オフセットはデータが Flume チャンネルに送信された後にのみコミットされます。

datahub.offsetCommitTimeout

30

オプション

オフセットのオートコミット間隔 (秒単位)。

datahub.sessionTimeout

60

オプション

ソース機能は協調消費を使用します。タイムアウト期間内にハートビートが送信されない場合、セッションは自動的に閉じられます。

serializer

-

必須

データ解析方法。現在、DELIMITED (区切り文字) のみがサポートされています。データの各フィールドは、DataHub スキーマの順序で行として書き込まれ、指定されたデリミタで区切られます。

serializer.delimiter

,

オプション

データフィールドの区切り文字。特殊文字を使用する場合は、"\t" のように二重引用符で囲みます。

serializer.charset

UTF-8

オプション

データ解析のエンコード形式。

ユースケース

Sink の例

例 1: DELIMITED シリアライザー

  1. テストデータを準備します。

    DELIMITED シリアライザーを使用する場合、各行はレコードとして扱われ、指定された区切り文字を使用して解析されます。次の例は、Flume を使用してバッチ CSV ファイルをほぼリアルタイムで DataHub にアップロードする方法を示しています。以下の内容を /temp/test.csv という名前のローカルファイルに保存します。

    0,YxCOHXcst1NlL5ebJM9YmvQ1f8oy8neb3obdeoS0,true,1254275.1144629316,1573206062763,1254275.1144637289
    0,YxCOHXcst1NlL5ebJM9YmvQ1f8oy8neb3obdeoS0,true,1254275.1144629316,1573206062763,1254275.1144637289
    1,hHVNjKW5DsRmVXjguwyVDjzjn60wUcOKos9Qym0V,false,1254275.1144637289,1573206062763,1254275.1144637289
    2,vnXOEuKF4Xdn5WnDCPbzPwTwDj3k1m3rlqc1vN2l,true,1254275.1144637289,1573206062763,1254275.1144637289
    3,t0AGT8HShzroBVM3vkP37fIahg2yDqZ5xWfwDFJs,false,1254275.1144637289,1573206062763,1254275.1144637289
    4,MKwZ1nczmCBp6whg1lQeFLZ6E628lXvFncUVcYWI,true,1254275.1144637289,1573206062763,1254275.1144637289
    5,bDPQJ656xvPGw1PPjhhTUZyLJGILkNnpqNLaELWV,false,1254275.1144637289,1573206062763,1254275.1144637289
    6,wWF7i4X8SXNhm4EfClQjQF4CUcYQgy3XnOSz0StX,true,1254275.1144637289,1573206062763,1254275.1144637289
    7,whUxTNREujMP6ZrAJlSVhCEKH1KH9XYJmOFXKbh8,false,1254275.1144637289,1573206062763,1254275.1144637289
    8,OYcS1WkGcbZFbPLKaqU5odlBf7rHDObkQJdBDrYZ,true,1254275.1144637289,1573206062763,1254275.1144637289

    テストデータに対応する DataHub スキーマは次のとおりです。

    フィールド名

    フィールドタイプ

    id

    BIGINT

    名前

    STRING

    性別

    ブーリアン

    給与

    DOUBLE

    my_time

    タイムスタンプ

    10 進数

    DECIMAL

  2. Flume ファイルを構成します。

    ${FLUME_HOME}/conf ディレクトリに、datahub_basic.conf という名前のファイルを作成し、以下の内容を追加します。この例では、Exec Source をデータソースとして使用します。他のソースの詳細については、公式 Flume ドキュメントをご参照ください。

    # A single-node Flume configuration for DataHub
    # Name the components on this agent
    a1.sources = r1
    a1.sinks = k1
    a1.channels = c1
    # Describe/configure the source
    a1.sources.r1.type = exec
    a1.sources.r1.command = cat /temp/test.csv
    # Describe the sink
    a1.sinks.k1.type = com.aliyun.datahub.flume.sink.DatahubSink
    a1.sinks.k1.datahub.accessId = {YOUR_ALIYUN_DATAHUB_ACCESS_ID}
    a1.sinks.k1.datahub.accessKey = {YOUR_ALIYUN_DATAHUB_ACCESS_KEY}
    a1.sinks.k1.datahub.endPoint = {YOUR_ALIYUN_DATAHUB_ENDPOINT}
    a1.sinks.k1.datahub.project = datahub_project_test
    a1.sinks.k1.datahub.topic = test_topic
    a1.sinks.k1.serializer = DELIMITED
    a1.sinks.k1.serializer.delimiter = ,
    a1.sinks.k1.serializer.fieldnames = id,name,gender,salary,my_time,decimal
    a1.sinks.k1.serializer.charset = UTF-8
    a1.sinks.k1.datahub.retryTimes = 5
    a1.sinks.k1.datahub.retryInterval = 5
    a1.sinks.k1.datahub.batchSize = 100
    a1.sinks.k1.datahub.batchTimeout = 5
    a1.sinks.k1.datahub.enablePb = true
    a1.sinks.k1.datahub.compressType = DEFLATE
    # Use a channel which buffers events in memory
    a1.channels.c1.type = memory
    a1.channels.c1.capacity = 10000
    a1.channels.c1.transactionCapacity = 10000
    # Bind the source and sink to the channel
    a1.sources.r1.channels = c1
    a1.sinks.k1.channel = c1
    説明

    ExecSource ソースは、イベントがチャンネルに配置されることを保証しないため、データ損失を引き起こす可能性があります。 例えば、`tail` コマンドがデータを取得する際に Flume チャンネルが満杯の場合、そのデータは失われます。 Spooling Directory Source または Taildir Source の使用を推奨します。 この例では、静的ファイル /temp/test.csv をデータソースとして使用します。 ファイルが動的に書き込まれるログファイルの場合、tail -F logFile コマンドを使用してリアルタイム収集を実行できます。

  3. Flume を起動します。

    `Dflume.root.logger=INFO,console` オプションは、ログをリアルタイムでコンソールに出力します。詳細情報を取得するには、DEBUG モードを使用します。次のコマンドを実行して Flume を起動し、CSV ファイルから DataHub にデータをインジェストします。

    $ cd ${FLUME_HOME}
    $ bin/flume-ng agent -n a1 -c conf -f conf/datahub_basic.conf -Dflume.root.logger=INFO,console

例 2: REGEX シリアライザー

  1. テストデータを準備します。

    REGEX シリアライザーを使用する場合、各行はレコードとして扱われ、指定された正規表現を使用して解析されます。レコードの異なる部分は、式内のグループで表されます。次の例は、Flume と正規表現を使用して、ログファイルから DataHub にデータをほぼリアルタイムでアップロードする方法を示しています。以下のテストデータを /temp/test.csv という名前のローカルファイルに保存します。

    1. [2019-11-12 15:20:08] 0,j4M6PhzL1DXVTQawdfk306N2KnCDxtR0KK1pke5O,true,1254409.5059812006,1573543208698,1254409.5059819978
    2. [2019-11-12 15:22:35] 0,mYLF8UzIYCCFUm1jYs9wzd2Hl6IMr2N7GPYXZSZy,true,1254409.5645912462,1573543355740,1254409.5645920434
    3. [2019-11-12 15:23:14] 0,MOemUZur37n4SGtdUQyMohgmM6cxZRBXjJ34HzqX,true,1254409.5799291395,1573543394219,1254409.579929538
    4. [2019-11-12 15:23:30] 0,EAFc1VTOvC9rYzPl9zJYa6cc8uJ089EaFd79B25i,true,1254409.5862723626,1573543410134,1254409.5862731598
    5. [2019-11-12 15:23:53] 0,zndVraA4GP7FP8p4CkQFsKJkxwtYK3zXjDdkhmRk,true,1254409.5956010541,1573543433538,1254409.5956018514
    6. [2019-11-12 15:24:00] 0,9YrjjoALEfyZm07J7OuNvDVNyspIzrbOOAGnZtHx,true,1254409.598201082,1573543440061,1254409.5982018793
    7. [2019-11-12 15:24:23] 0,mWsFgFlUnXKQQR6RpbAYDF9OhGYgU8mljvGCtZ26,true,1254409.6073950487,1573543463126,1254409.607395447
    8. [2019-11-12 15:26:51] 0,5pZRRzkW3WDLdYLOklNgTLFX0Q0uywZ8jhw7RYfI,true,1254409.666525653,1573543611475,1254409.6665264503
    9. [2019-11-12 15:29:11] 0,hVgGQrXpBtTJm6sovVK4YGjfNMdQ3z9pQHxD5Iqd,true,1254409.7222845491,1573543751364,1254409.7222853464
    10. [2019-11-12 15:29:52] 0,7wQOQmxoaEl6Cxl1OSo6cr8MAc1AdJWJQaTPT5xs,true,1254409.7387664048,1573543792714,1254409.738767202
    11. [2019-11-12 15:30:30] 0,a3Th5Q6a8Vy2h1zfWLEP7MdPhbKyTY3a4AfcOJs2,true,1254409.7538966285,1573543830673,1254409.7538974257
    12. [2019-11-12 15:34:54] 0,d0yQAugqJ8M8OtmVQYMTYR8hi3uuX5WsH9VQRBpP,true,1254409.8589555968,1573544094247,1254409.8589563938

    テストデータに対応する DataHub スキーマは次のとおりです。

    フィールド名

    フィールドタイプ

    id

    BIGINT

    name

    STRING

    gender

    BOOLEAN

    salary

    DOUBLE

    my_time

    TIMESTAMP

    decimal

    DECIMAL

  2. Flume ファイルを構成します。

    ${FLUME_HOME}/conf ディレクトリに datahub_basic.conf という名前でファイルを作成し、次の内容を追加します。 この例では、データソースとして Exec Source を使用します。 他のソースの詳細については、「Flume の公式ドキュメント」をご参照ください。

    # A single-node Flume configuration for DataHub
    # Name the components on this agent
    a1.sources = r1
    a1.sinks = k1
    a1.channels = c1
    # Describe/configure the source
    a1.sources.r1.type = exec
    a1.sources.r1.command = cat /temp/test.csv
    # Describe the sink
    a1.sinks.k1.type = com.aliyun.datahub.flume.sink.DatahubSink
    a1.sinks.k1.datahub.accessId = {YOUR_ALIYUN_DATAHUB_ACCESS_ID}
    a1.sinks.k1.datahub.accessKey = {YOUR_ALIYUN_DATAHUB_ACCESS_KEY}
    a1.sinks.k1.datahub.endPoint = {YOUR_ALIYUN_DATAHUB_ENDPOINT}
    a1.sinks.k1.datahub.project = datahub_project_test
    a1.sinks.k1.datahub.topic = test_topic
    a1.sinks.k1.serializer = REGEX
    a1.sinks.k1.serializer.regex = \\[\\d{4}-\\d{2}-\\d{2} \\d{2}:\\d{2}:\\d{2}\\] (\\d+),(\\S+),([a-z]+),([-+]?[0-9]*\\.?[0-9]*),(\\d+),([-+]?[0-9]*\\.?[0-9]*)
    a1.sinks.k1.serializer.fieldnames = id,name,gender,salary,my_time,decimal
    a1.sinks.k1.serializer.charset = UTF-8
    a1.sinks.k1.datahub.retryTimes = 5
    a1.sinks.k1.datahub.retryInterval = 5
    a1.sinks.k1.datahub.batchSize = 100
    a1.sinks.k1.datahub.batchTimeout = 5
    # Use a channel which buffers events in memory
    a1.channels.c1.type = memory
    a1.channels.c1.capacity = 10000
    a1.channels.c1.transactionCapacity = 10000
    # Bind the source and sink to the channel
    a1.sources.r1.channels = c1
    a1.sinks.k1.channel = c1
    説明

    ExecSource ソースは、イベントがチャンネルに配置されることを保証しないため、データ損失を引き起こす可能性があります。 例えば、`tail` コマンドがデータを取得するときに Flume チャンネルがいっぱいの場合、そのデータは失われます。 そのため、Spooling Directory Source または Taildir Source の使用をお勧めします。 この例では、静的ファイル /temp/test.csv をデータソースとして使用します。 ファイルが動的に書き込まれるログファイルの場合、tail -F logFile コマンドを使用してリアルタイム収集を実行できます。

  3. Flume を起動します。

    Dflume.root.logger=INFO,console オプションは、ログをコンソールにリアルタイムで出力します。より詳細な情報を取得するには、DEBUG モードを使用します。次のコマンドを実行して、Flume を起動し、CSV ファイルから DataHub へデータをインジェストします:

    $ cd ${FLUME_HOME}
    $ bin/flume-ng agent -n a1 -c conf -f conf/datahub_basic.conf -Dflume.root.logger=INFO,console

例 3: Flume Taildir Source

Flume の `exec` ソースはデータ損失を引き起こす可能性があります。したがって、本番環境での使用は推奨されません。ローカルログを収集するには、`Taildir Source` または `Spooling Directory Source` を使用できます。次の例は、Taildir を使用してログファイルを収集する方法を示しています。`Taildir Source` は、指定されたファイルのグループをモニターし、各ファイルに追加される新しい行をほぼリアルタイムで読み取ることができます。新しい行が書き込まれている場合、このソースは書き込み操作が完了するまで読み取りを再試行します。`Taildir Source` は、各ファイルの読み取り位置を `positionFile` という名前の JSON ファイルに格納します。ソースイベントがチャンネルに配置できない場合、読み取り位置は更新されません。これにより、`Taildir Source` は信頼性があります。

  1. テストデータを準備します。

    すべてのログは、次のフォーマットでファイルに追加されます。ログファイルは *.log という命名フォーマットを使用します。

    0,YxCOHXcst1NlL5ebJM9YmvQ1f8oy8neb3obdeoS0,true,1254275.1144629316,1573206062763,1254275.1144637289

    テストデータに対応する DataHub スキーマは次のとおりです。

    フィールド名

    フィールドタイプ

    id

    BIGINT

    name

    STRING

    gender

    BOOLEAN

    salary

    DOUBLE

    my_time

    TIMESTAMP

    decimal

    DECIMAL

  2. Flume 構成ファイルを構成します。

    ${FLUME_HOME}/conf ディレクトリに、 datahub_basic.conf という名前のファイルを作成し、ファイルに次の内容を追加します。

    # A single-node Flume configuration for DataHub
    # Name the components on this agent
    a1.sources = r1
    a1.sinks = k1
    a1.channels = c1
    # Describe/configure the source
    a1.sources.r1.type = TAILDIR
    a1.sources.r1.positionFile = /temp/taildir_position.json
    a1.sources.r1.filegroups = f1
    a1.sources.r1.filegroups.f1 = /temp/.*log
    # Describe the sink
    a1.sinks.k1.type = com.aliyun.datahub.flume.sink.DatahubSink
    a1.sinks.k1.datahub.accessId = {YOUR_ALIYUN_DATAHUB_ACCESS_ID}
    a1.sinks.k1.datahub.accessKey = {YOUR_ALIYUN_DATAHUB_ACCESS_KEY}
    a1.sinks.k1.datahub.endPoint = {YOUR_ALIYUN_DATAHUB_ENDPOINT}
    a1.sinks.k1.datahub.project = datahub_project_test
    a1.sinks.k1.datahub.topic = test_topic
    a1.sinks.k1.serializer = DELIMITED
    a1.sinks.k1.serializer.delimiter = ,
    a1.sinks.k1.serializer.fieldnames = id,name,gender,salary,my_time,decimal
    a1.sinks.k1.serializer.charset = UTF-8
    a1.sinks.k1.datahub.retryTimes = 5
    a1.sinks.k1.datahub.retryInterval = 5
    a1.sinks.k1.datahub.batchSize = 100
    a1.sinks.k1.datahub.batchTimeout = 5
    a1.sinks.k1.datahub.enablePb = true
    a1.sinks.k1.datahub.compressType = DEFLATE
    # Use a channel which buffers events in memory
    a1.channels.c1.type = memory
    a1.channels.c1.capacity = 10000
    a1.channels.c1.transactionCapacity = 10000
    # Bind the source and sink to the channel
    a1.sources.r1.channels = c1
    a1.sinks.k1.channel = c1
  3. Flume を起動します。

    Dflume.root.logger=INFO,console オプションは、ログをリアルタイムでコンソールに出力します。詳細情報を取得するには、DEBUG モードを使用します。次のコマンドを実行して Flume を起動し、CSV ファイルから DataHub にデータを取り込みます:

    1. $ cd ${FLUME_HOME}
    2. $ bin/flume-ng agent -n a1 -c conf -f conf/datahub_basic.conf -Dflume.root.logger=INFO,console

例 4: JSON シリアライザー

JSON シリアライザーを使用する場合、各行はレコードとして扱われます。JSON オブジェクトの最初のレイヤーのみが解析されます。ネストされたコンテンツは文字列として扱われます。構成された serializer.fieldnames にトップレベル名が存在する場合、その値は対応する列に追加されます。次の例は、Flume と JSON 解析を使用して、ログファイルから DataHub にデータをほぼリアルタイムでアップロードする方法を示しています。

  1. テストデータを準備します。

    以下の内容をローカルファイル /temp/test.json に保存します。同期するデータは、日付の後に表示される詳細情報です。

    {"my_time":1573206062763,"gender":true,"name":"YxCOHXcst1NlL5ebJM9YmvQ1f8oy8neb3obdeoS0","id":0,"salary":1254275.1144629316,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":true,"name":"YxCOHXcst1NlL5ebJM9YmvQ1f8oy8neb3obdeoS0","id":0,"salary":1254275.1144629316,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":false,"name":"hHVNjKW5DsRmVXjguwyVDjzjn60wUcOKos9Qym0V","id":1,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":true,"name":"vnXOEuKF4Xdn5WnDCPbzPwTwDj3k1m3rlqc1vN2l","id":2,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":false,"name":"t0AGT8HShzroBVM3vkP37fIahg2yDqZ5xWfwDFJs","id":3,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":true,"name":"MKwZ1nczmCBp6whg1lQeFLZ6E628lXvFncUVcYWI","id":4,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":false,"name":"bDPQJ656xvPGw1PPjhhTUZyLJGILkNnpqNLaELWV","id":5,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":true,"name":"wWF7i4X8SXNhm4EfClQjQF4CUcYQgy3XnOSz0StX","id":6,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"my_time":1573206062763,"gender":false,"name":"whUxTNREujMP6ZrAJlSVhCEKH1KH9XYJmOFXKbh8","id":7,"salary":1254275.1144637289,"decimal":1254275.1144637289}
    {"gender":true,"name":{"a":"OYcS1WkGcbZFbPLKaqU5odlBf7rHDObkQJdBDrYZ"},"id":8,"salary":1254275.1144637289,"decimal":1254275.1144637289}

    テストデータに対応する DataHub スキーマは次のとおりです。

    フィールド名

    フィールドタイプ

    id

    BIGINT

    name

    STRING

    gender

    BOOLEAN

    salary

    DOUBLE

    my_time

    TIMESTAMP

    decimal

    DECIMAL

  2. Flume ファイルを構成します。

    ${FLUME_HOME}/conf ディレクトリに、datahub_basic.conf という名前のファイルを作成し、次の内容を追加します。この例では、データソースとして Exec Source を使用します。他のソースに関する詳細については、「Flume の公式ドキュメント」をご参照ください。

    # A single-node Flume configuration for DataHub
    # Name the components on this agent
    a1.sources = r1
    a1.sinks = k1
    a1.channels = c1
    # Describe/configure the source
    a1.sources.r1.type = exec
    a1.sources.r1.command = cat /temp/test.json
    # Describe the sink
    a1.sinks.k1.type = com.aliyun.datahub.flume.sink.DatahubSink
    a1.sinks.k1.datahub.accessId = {YOUR_ALIYUN_DATAHUB_ACCESS_ID}
    a1.sinks.k1.datahub.accessKey = {YOUR_ALIYUN_DATAHUB_ACCESS_KEY}
    a1.sinks.k1.datahub.endPoint = {YOUR_ALIYUN_DATAHUB_ENDPOINT}
    a1.sinks.k1.datahub.project = datahub_project_test
    a1.sinks.k1.datahub.topic = test_topic
    a1.sinks.k1.serializer = JSON
    a1.sinks.k1.serializer.fieldnames = id,name,gender,salary,my_time,decimal
    a1.sinks.k1.serializer.charset = UTF-8
    a1.sinks.k1.datahub.retryTimes = 5
    a1.sinks.k1.datahub.retryInterval = 5
    a1.sinks.k1.datahub.batchSize = 100
    a1.sinks.k1.datahub.batchTimeout = 5
    # Use a channel which buffers events in memory
    a1.channels.c1.type = memory
    a1.channels.c1.capacity = 10000
    a1.channels.c1.transactionCapacity = 10000
    # Bind the source and sink to the channel
    a1.sources.r1.channels = c1
    a1.sinks.k1.channel = c1
  3. Flume を起動します。

    Dflume.root.logger=INFO,console オプションは、ログをリアルタイムでコンソールに出力します。詳細情報を取得するには、DEBUG モードを使用します。次のコマンドを実行して Flume を起動し、CSV ファイルから DataHub にデータを取り込みます:

    $ cd ${FLUME_HOME}
    $ bin/flume-ng agent -n a1 -c conf -f conf/datahub_basic.conf -Dflume.root.logger=INFO,console

Source の例

DataHub から他のシステムへのデータ読み取り

DataHub-Flume ソースを使用して、DataHub からデータを読み取り、別のシステムに移動できます。このトピックでは、コンソールに直接出力するロガーシンクを例として、DataHub-Flume ソースの使用方法を示します。

  1. 以下はトピックスキーマの例です。

    フィールド名

    フィールドタイプ

    id

    BIGINT

    name

    STRING

    gender

    BOOLEAN

    salary

    DOUBLE

    my_time

    TIMESTAMP

    decimal

    DECIMAL

  2. Flume ファイルを構成します。

    ${FLUME_HOME}/conf ディレクトリで、datahub_source.confという名前のファイルを作成し、ファイルに以下の内容を追加します。

     # A single-node Flume configuration for DataHub
     # Name the components on this agent
     a1.sources = r1
     a1.sinks = k1
     a1.channels = c1
    
     # Describe/configure the source
     a1.sources.r1.type = com.aliyun.datahub.flume.sink.DatahubSource
     a1.sources.r1.datahub.endPoint = {YOUR_ALIYUN_DATAHUB_ENDPOINT}
     a1.sources.r1.datahub.accessId = {YOUR_ALIYUN_DATAHUB_ACCESS_ID}
     a1.sources.r1.datahub.accessKey = {YOUR_ALIYUN_DATAHUB_ACCESS_KEY}
     a1.sources.r1.datahub.project = datahub_test
     a1.sources.r1.datahub.topic = test_flume
     a1.sources.r1.datahub.subId = {YOUR_ALIYUN_DATAHUB_SUB_ID}
     a1.sources.r1.serializer = DELIMITED
     a1.sources.r1.serializer.delimiter = ,
     a1.sources.r1.serializer.charset = UTF-8
     a1.sources.r1.datahub.retryTimes = 3
     a1.sources.r1.datahub.batchSize = 1000
     a1.sources.r1.datahub.batchTimeout = 5
     a1.sources.r1.datahub.enablePb = false
    
     # Describe the sink
     a1.sinks.k1.type = logger
    
     # Use a channel which buffers events in memory
     a1.channels.c1.type = memory
     a1.channels.c1.capacity = 10000
     a1.channels.c1.transactionCapacity = 10000
    
     # Bind the source and sink to the channel
     a1.sources.r1.channels = c1
     a1.sinks.k1.channel = c1
  3. Flume を起動します。

    $ cd ${FLUME_HOME}
    $ bin/flume-ng agent -n a1 -c conf -f conf/datahub_source.conf -Dflume.root.logger=INFO,console

Flume メトリック

DataHub-Flume は、Flume の組み込みカウンターモニターをサポートしており、これを使用して Flume プラグインの実行状態を監視できます。DataHub-Flume プラグインのシンクとソースは、メトリック情報を表示できます。次の表は、DataHub 関連パラメーターについて説明しています。他のパラメーターの詳細については、公式 Flume ドキュメントをご参照ください。

DatahubSink

名前

説明

BatchEmptyCount

DataHub に書き込むデータがない状態でバッチがタイムアウトした回数。

BatchCompleteCount

正常に処理されたバッチの数。これは、すべてのデータが正常に書き込まれた場合のみを含みます。

EventDrainAttemptCount

DataHub への書き込みが試行されたレコード数 (正常に解析されたレコード数)。

BatchUnderflowCount

DataHub に正常に書き込まれたデータ量が、書き込む必要があったデータ量よりも少なかった回数。これは、データ解析が完了しても、DataHub への書き込みが部分的または完全に失敗した場合に発生します。

EventDrainSuccessCount

DataHub に正常に書き込まれたデータ量。

DatahubSource

名前

説明

EventReceivedCount

ソースが DataHub から受信したレコード数。

EventAcceptedCount

ソースが DataHub からチャンネルに正常に書き込んだレコード数。

Flume モニタリング

Flume は複数のモニタリング方法を提供します。このトピックでは、HTTP モニタリングを例として、Flume のモニタリングツールの使用方法を示します。HTTP モニタリングを使用するには、Flume プラグインを起動するときに 2 つのパラメーターを追加します: -Dflume.monitoring.type=http -Dflume.monitoring.port=1234。`type` パラメーターはモニタリング方法を指定し、`port` パラメーターはポート番号を指定します。以下に例を示します。

bin/flume-ng agent -n a1 -c conf -f conf/datahub_basic.conf -Dflume.root.logger=INFO,console -Dflume.monitoring.type=http -Dflume.monitoring.port=1234

プラグインの起動後、https://ip:1234/metrics の Web UI でメトリックを表示できます。

説明

モニタリング方法の詳細については、公式 Flume ドキュメントをご参照ください。

よくある質問

Flume が起動に失敗し、エラーを報告: org.apache.flume.ChannelFullException: Space for commit to queue couldn’t be acquired. Sinks are likely not keeping up with sources, or the buffer size is too tight

Flume のデフォルトのヒープメモリは 20 MB です。`batchSize` パラメーターを大きな値に設定すると、Flume が使用するヒープメモリが 20 MB を超える可能性があります。

解決策 1: `batchSize` の値を減らします。

解決策 2: Flume の最大ヒープメモリを増やします。

  • $ vim bin/flume-ng

  • JAVA_OPTS="-Xmx20m" ==> JAVA_OPTS="-Xmx1024m"

DataHub-Flume プラグインは JSON フォーマットをサポートしていますか?

いいえ、サポートしていません。ただし、カスタム正規表現を使用してデータを解析するか、DataHub-Flume プラグインのコードを変更して JSONEvent のサポートを追加できます。

DataHub-Flume プラグインは BLOB Topic をサポートしていますか?

DataHub-Flume プラグインは現在、Tuple トピックのみをサポートしており、BLOB Topic はサポートしていません。

Flume がエラーを報告: org.apache.flume.ChannelException: Put queue for MemoryTransaction of capacity 1 full, consider committing more frequently, increasing capacity or increasing thread count

このエラーは、チャンネルが満杯で、ソースがチャンネルにデータを書き込めなかったために発生します。この問題を解決するには、構成ファイルでチャンネル容量を増やし、DataHub ソースの `batchSize` を減らすことができます。

古いバージョンの Flume を使用するとエラーが発生し、JAR パッケージの競合により起動に失敗する可能性

  • シナリオ: Flume 1.6 を使用すると、起動が失敗し、次のエラーが報告される場合があります: java.lang.NoSuchMethodError:com.fasterxml.jackson.databind.ObjectMapper.readerFor(Lcom/fasterxml/jackson/databind/JavaType;)Lcom/fasterxml/jackson/databind/ObjectReader;。このエラーは、新しいプラグインが依存する JAR パッケージが Flume が依存するバージョンと一致しないために発生します。Flume の古い JAR パッケージを使用すると、新しいメソッドが見つかりません。

  • 解決策: ${FLUME_HOME}/lib ディレクトリから次の 3 つの JAR パッケージを削除します。

    • jackson-annotations-2.3.0.jar

    • jackson-databind-2.3.1.jar

    • jackson-annotations-2.3.0.jar

Flume でのデータインジェスト中に空の文字列が自動的に null に変換される

Flume プラグインのバージョン 2.0.2 では、空でない文字列はトリミングされ、空の文字列は null に変換されます。この問題はバージョン 2.0.3 で修正されています。バージョン 2.0.3 では、空の文字列は DataHub に空の文字列として書き込まれます。

起動がエラーで失敗: Cannot invoke "com.google.common.cache.LoadingCache.get(Object)" because"com.aliyun.datahub.client.impl.batch.avro.AvroSchemaCache.schemaCache" is null]

Flume の `lib` フォルダーから `guava` および `zstd` JAR ファイルを削除し、Flume を再起動します。