ブローカーロードは、HDFS、Baidu Object Storage (BOS)、または Andrew File System (AFS) から Doris へ大量のデータをロードするための非同期インポートメソッドです。Spark コンピューティングリソースが利用できない場合で、ブローカーアクセス可能なファイルシステムに保存されている数十から数百 GB のデータをインポートする必要がある場合に使用します。
ブローカーロードは、単一のインポートジョブに対する唯一の非同期インポートメソッドです。Spark リソースが利用可能な場合は、代わりに Spark Load を使用してください。これは、大容量の既存データ移行において Doris クラスターリソースの使用量を削減します。
仕組み
インポートジョブを送信すると、Frontend (FE) は実行計画を生成し、BE の数とソースファイルサイズに基づいて利用可能な Backends (BE) に分散します。各 BE はブローカーからデータをプルし、変換してシステムにロードします。すべての BE が完了すると、FE はジョブが成功したかどうかを判断します。
| 1. User submits a Broker Load job
v
+----+----+
| |
| FE |
| |
+----+----+
|
| 2. Each BE runs extract, transform, and load
+--------------------------+
| | |
+---v---+ +--v----+ +---v---+
| | | | | |
| BE | | BE | | BE |
| | | | | |
+---+-^-+ +---+-^-+ +--+-^--+
| | | | | |
| | | | | | 3. Each BE pulls data from a broker
+---v-+-+ +---v-+-+ +--v-+--+
| | | | | |
|Broker | |Broker | |Broker |
| | | | | |
+---+-^-+ +---+-^-+ +---+-^-+
| | | | | |
+---v-+-----------v-+----------v-+-+
| HDFS/BOS/AFS Cluster |
| |
+----------------------------------+前提条件
開始する前に、以下を確認してください。
少なくとも 1 つの BE を持つ実行中の Doris クラスター
ブローカーアクセス可能なファイルシステム (HDFS、BOS、または AFS) に保存されているソースデータ
インポート用に作成された Doris ターゲットテーブル
インポートジョブの送信
パーティション化された Hive テーブルからのデータインポート
この例では、デフォルトの CSV 形式を使用して、パーティション化された Hive テーブルからデータをインポートします。
dayパーティションフィールドを持つ Hive テーブルを作成します。-- Default data format (CSV), partition field: day CREATE TABLE `ods_demo_detail`( `id` string, `store_id` string, `company_id` string, `tower_id` string, `commodity_id` string, `commodity_name` string, `commodity_price` double, `member_price` double, `cost_price` double, `unit` string, `quantity` double, `actual_price` double ) PARTITIONED BY (day string) row format delimited fields terminated by ',' lines terminated by '\n'Hive テーブルにデータをロードします。
load data local inpath '/opt/custorm' into table ods_demo_detail;ターゲット Doris テーブルを作成します。
CREATE TABLE `doris_ods_test_detail` ( `rq` date NULL, `id` varchar(32) NOT NULL, `store_id` varchar(32) NULL, `company_id` varchar(32) NULL, `tower_id` varchar(32) NULL, `commodity_id` varchar(32) NULL, `commodity_name` varchar(500) NULL, `commodity_price` decimal(10, 2) NULL, `member_price` decimal(10, 2) NULL, `cost_price` decimal(10, 2) NULL, `unit` varchar(50) NULL, `quantity` int(11) NULL, `actual_price` decimal(10, 2) NULL ) ENGINE=OLAP UNIQUE KEY(`rq`, `id`, `store_id`) PARTITION BY RANGE(`rq`) ( PARTITION P_202204 VALUES [('2022-04-01'), ('2022-05-01'))) DISTRIBUTED BY HASH(`store_id`) BUCKETS 1 PROPERTIES ( "replication_allocation" = "tag.location.default: 3", "dynamic_partition.enable" = "true", "dynamic_partition.time_unit" = "MONTH", "dynamic_partition.start" = "-2147483648", "dynamic_partition.end" = "2", "dynamic_partition.prefix" = "P_", "dynamic_partition.buckets" = "1", "in_memory" = "false", "storage_format" = "V2" );インポートジョブを送信します。
SET句は、Hive パーティションフィールドdayをstr_to_dateを使用して Doris の日付列rqにマッピングします。LOAD LABEL broker_load_2022_03_23 ( DATA INFILE("hdfs://192.168.**.**:8020/user/hive/warehouse/ods.db/ods_demo_detail/*/*") INTO TABLE doris_ods_test_detail COLUMNS TERMINATED BY "," (id,store_id,company_id,tower_id,commodity_id,commodity_name,commodity_price,member_price,cost_price,unit,quantity,actual_price) COLUMNS FROM PATH AS (`day`) SET (rq = str_to_date(`day`,'%Y-%m-%d'),id=id,store_id=store_id,company_id=company_id,tower_id=tower_id,commodity_id=commodity_id,commodity_name=commodity_name,commodity_price=commodity_price,member_price=member_price,cost_price=cost_price,unit=unit,quantity=quantity,actual_price=actual_price) ) WITH BROKER "broker_name_1" ( "username" = "hdfs", "password" = "" ) PROPERTIES ( "timeout"="1200", "max_filter_ratio"="0.1" );
ORC 形式のパーティション化された Hive テーブルからのデータインポート
この例では、Optimized Row Columnar (ORC) 形式で保存されている Hive テーブルからインポートします。CSV の例との主な違いは、FORMAT AS "orc" 句です。
ORC 形式でパーティション化された Hive テーブルを作成します。
-- Data format: ORC, partition field: day CREATE TABLE `ods_demo_orc_detail`( `id` string, `store_id` string, `company_id` string, `tower_id` string, `commodity_id` string, `commodity_name` string, `commodity_price` double, `member_price` double, `cost_price` double, `unit` string, `quantity` double, `actual_price` double ) PARTITIONED BY (day string) row format delimited fields terminated by ',' lines terminated by '\n' STORED AS ORC前の例と同じスキーマを使用して、ターゲット Doris テーブルを作成します。
CREATE TABLE `doris_ods_test_detail` ( `rq` date NULL, `id` varchar(32) NOT NULL, `store_id` varchar(32) NULL, `company_id` varchar(32) NULL, `tower_id` varchar(32) NULL, `commodity_id` varchar(32) NULL, `commodity_name` varchar(500) NULL, `commodity_price` decimal(10, 2) NULL, `member_price` decimal(10, 2) NULL, `cost_price` decimal(10, 2) NULL, `unit` varchar(50) NULL, `quantity` int(11) NULL, `actual_price` decimal(10, 2) NULL ) ENGINE=OLAP UNIQUE KEY(`rq`, `id`, `store_id`) PARTITION BY RANGE(`rq`) ( PARTITION P_202204 VALUES [('2022-04-01'), ('2022-05-01'))) DISTRIBUTED BY HASH(`store_id`) BUCKETS 1 PROPERTIES ( "replication_allocation" = "tag.location.default: 3", "dynamic_partition.enable" = "true", "dynamic_partition.time_unit" = "MONTH", "dynamic_partition.start" = "-2147483648", "dynamic_partition.end" = "2", "dynamic_partition.prefix" = "P_", "dynamic_partition.buckets" = "1", "in_memory" = "false", "storage_format" = "V2" );FORMAT AS "orc"を指定してインポートジョブを送信します。LOAD LABEL dish_2022_03_23 ( DATA INFILE("hdfs://10.220.**.**:8020/user/hive/warehouse/ods.db/ods_demo_orc_detail/*/*") INTO TABLE doris_ods_test_detail COLUMNS TERMINATED BY "," FORMAT AS "orc" (id,store_id,company_id,tower_id,commodity_id,commodity_name,commodity_price,member_price,cost_price,unit,quantity,actual_price) COLUMNS FROM PATH AS (`day`) SET (rq = str_to_date(`day`,'%Y-%m-%d'),id=id,store_id=store_id,company_id=company_id,tower_id=tower_id,commodity_id=commodity_id,commodity_name=commodity_name,commodity_price=commodity_price,member_price=member_price,cost_price=cost_price,unit=unit,quantity=quantity,actual_price=actual_price) ) WITH BROKER "broker_name_1" ( "username" = "hdfs", "password" = "" ) PROPERTIES ( "timeout"="1200", "max_filter_ratio"="0.1" );
HDFS からの直接データインポート
Hive テーブルを経由せずに HDFS からタブ区切りテキストファイルを直接インポートするには、WITH BROKER 句の代わりに WITH HDFS 句を使用します。
LOAD LABEL demo.label_20220402
(
DATA INFILE("hdfs://10.220.**.**:8020/tmp/test_hdfs.txt")
INTO TABLE `ods_dish_detail_test`
COLUMNS TERMINATED BY "\t" (id,store_id,company_id,tower_id,commodity_id,commodity_name,commodity_price,member_price,cost_price,unit,quantity,actual_price)
)
WITH HDFS (
"fs.defaultFS"="hdfs://10.220.**.**:8020",
"hadoop.username"="root"
)
PROPERTIES
(
"timeout"="1200",
"max_filter_ratio"="0.1"
);インポートジョブステータスの確認
次のステートメントを実行して、最新のインポートジョブを表示します。
show load order by createtime desc limit 1\G;出力例:
*************************** 1. row ***************************
JobId: 4132****
Label: broker_load_2022_03_23
State: FINISHED
Progress: ETL:100%; LOAD:100%
Type: BROKER
EtlInfo: unselected.rows=0; dpp.abnorm.ALL=0; dpp.norm.ALL=27
TaskInfo: cluster:N/A; timeout(s):1200; max_filter_ratio:0.1
ErrorMsg: NULL
CreateTime: 2022-04-01 18:59:06
EtlStartTime: 2022-04-01 18:59:11
EtlFinishTime: 2022-04-01 18:59:11
LoadStartTime: 2022-04-01 18:59:11
LoadFinishTime: 2022-04-01 18:59:11
URL: NULL
JobDetails: {"Unfinished backends":{"5072bde59b74b65-8d2c0ee5b029****":[]},"ScannedRows":27,"TaskNumber":1,"All backends":{"5072bde59b74b65-8d2c0ee5b029****":[36728051]},"FileNumber":1,"FileSize":5540}
1 row in set (0.01 sec)出力の主要フィールド:
| フィールド | 説明 |
|---|---|
State | ジョブステータス:PENDING、LOADING、FINISHED、または CANCELLED |
Progress | ETL およびロードの進捗率 |
EtlInfo | 行数:dpp.norm.ALL = 正常にロードされた行、dpp.abnorm.ALL = フィルタリングされた行 |
TaskInfo | タイムアウトと max_filter_ratio の有効な設定 |
ErrorMsg | ジョブが失敗した場合のエラーメッセージ。成功した場合は NULL |
インポートジョブのキャンセル
ジョブの状態が PENDING または LOADING の場合にのみジョブをキャンセルします。CANCEL LOAD ステートメントでジョブラベルを指定します。
CANCEL LOAD FROM demo WHERE LABEL = "broker_load_2022_03_23";システム構成
FE 構成
fe.conf のこれらのパラメーターは、クラスター内のすべてのブローカーロードジョブに適用され、FE が BE 間で作業をどのように分散するかを制御します。
| パラメーター | デフォルト | 単位 | 説明 |
|---|---|---|---|
min_bytes_per_broker_scanner | 64 MB | bytes | 単一の BE がジョブごとに処理する最小データ量 |
max_bytes_per_broker_scanner | 3 GB | bytes | 単一の BE がジョブごとに処理する最大データ量 |
max_broker_concurrency | 10 | — | インポートジョブごとの最大同時実行タスク数 |
desired_max_waiting_jobs | 100 | — | クラスター内の PENDING + LOADING ジョブの最大数。このしきい値を超えると新しいジョブは拒否されます。 |
async_pending_load_task_pool_size | 10 | — | 同時に実行される保留中タスクの最大数。LOADING 状態に入るジョブ数を制限します。 |
async_loading_load_task_pool_size | — | — | 同時に実行されるロード中タスクの最大数。async_pending_load_task_pool_size |
FE は、次の数式を使用して各インポートジョブの同時実行数を計算します。
Concurrency = Math.min(
Source file size / min_bytes_per_broker_scanner,
max_broker_concurrency,
Number of BEs
)
Data per BE = Source file size / Concurrency単一ジョブでインポート可能な最大データは、max_bytes_per_broker_scanner x Number of BEs と等しくなります。この制限を超えるデータをインポートするには、max_bytes_per_broker_scanner を増やします。
各ブローカーロードジョブは、1 つの保留中タスクと 1 つ以上のロード中タスクで構成されます。ロード中タスクの数は、LOAD LABELステートメントのDATA INFILE句の数と等しくなります。
ブローカーパラメーター
ブローカーごとに、リモートストレージにアクセスするために異なる認証パラメーターが必要です。必要な接続プロパティについては、ブローカードキュメントをご参照ください。
ベストプラクティス
ジョブごとのデータ量選択
以下のしきい値は、単一 BE クラスターを前提としています。マルチ BE クラスターの場合、BE の数で乗算します (例: 3 BE クラスターでは、3 GB 以下のしきい値が 9 GB 以下に引き上げられます)。
| データサイズ | アクション |
|---|---|
| <= 3 GB | デフォルト設定で直接送信 |
| > 3 GB | 送信する前に FE 構成を調整 (以下を参照) |
| > 500 GB | 複数のファイルに分割してバッチでインポート |
大容量ファイルの構成 (> 3 GB)
3 GB を超えるファイルの場合、ジョブを送信する前に次の手順に従います。
Step 1: fe.conf で同時実行数と BE ごとのデータ制限を設定します。
max_broker_concurrency = <number of BEs>
max_bytes_per_broker_scanner >= <source file size> / max_broker_concurrency例: 100 GB ファイル、10 BE。
max_broker_concurrency = 10
max_bytes_per_broker_scanner >= 10 GB (100 GB / 10)これらの設定により、すべての 10 BE がジョブを並行して処理し、それぞれ 10 GB を処理します。
これらの FE 設定は、現在のジョブだけでなく、クラスター内のすべてのブローカーロードジョブに適用されます。
Step 2: ジョブのタイムアウトを設定します。
クラスターのインポート速度を使用してタイムアウトを計算します。控えめな下限として 10 MB/秒 を使用します。
Data per BE / Slowest import speed of your cluster (MB/s) >= Timeout >= Data per BE / 10 MB/s例: BE ごとに 10 GB、タイムアウト >= 1,000 秒。
ジョブを送信する際に、PROPERTIES ブロックでタイムアウトを設定します。
Step 3: 4 時間以上かかるファイルを分割します。
デフォルトの最大タイムアウトは 4 時間です。計算がこれを超過する場合でも、最大タイムアウトを増やさず、代わりにソースファイルを分割してください。4 時間以上実行された失敗したジョブは、再試行にも同じ時間がかかります。
次の数式を使用して、4 時間以内に収まる最大ファイルサイズを見つけます。
Max data per batch = 14,400 s x 10 MB/s x Number of BEs例: 10 BE、バッチあたりの最大値 = 14,400 秒 x 10 MB/秒 x 10 = 1,440 GB。
実際には、インポート速度が 10 MB/秒に達することはめったにありません。500 GB を超えるファイルは、より小さなバッチに分割してください。
ジョブスケジューリングの管理
desired_max_waiting_jobs パラメーターは、PENDING または LOADING 状態にあるジョブ数を同時に制限します (デフォルト: 100)。このしきい値を超えて送信された新しいジョブは拒否されます。
async_pending_load_task_pool_size パラメーター (デフォルト: 10) は、LOADING 状態にアクティブに入るジョブ数を制限します。たとえば、100 個のジョブを送信した場合、同時に実行されるのは 10 個のみで、残りは PENDING 状態で待機します。