バッチ同期ノードでフィルター条件を使用すると、完全データまたは増分データのいずれかを同期できます。フィルター条件を使用すると、Data Integration は指定された基準を満たすデータのみを同期します。また、フィルター条件とスケジューリングパラメーターを組み合わせることで、ノードの実行時刻に基づいてデータを動的にフィルタリングし、増分同期を実現できます。本トピックでは、増分同期用のバッチ同期ノードを設定する方法について説明します。
注意事項
Hbase や Tablestore (OTSStream) など、一部のデータソースでは増分同期がサポートされていません。特定のデータソースで増分同期がサポートされているかどうかを確認するには、対応する Reader プラグインのドキュメントをご参照ください。
増分同期に必要なパラメーターは、Reader プラグインによって異なります。詳細については、特定のプラグインのドキュメントと サポートされているデータソースとプラグイン をご参照ください。例:
Reader プラグイン
必須パラメーター
サポートされている構文
where
説明ウィザードモードでは、これはフィルター条件パラメーターです。
データベース構文
説明このパラメーターをスケジューリングパラメーターと組み合わせて使用すると、毎日指定された時間範囲のデータを読み取ることができます。
query
説明ウィザードモードでは、これは検索条件パラメーターです。
データベース構文に類似
説明このパラメーターをスケジューリングパラメーターと組み合わせて使用すると、毎日指定された時間範囲のデータを読み取ることができます。
Object
オブジェクトパスを指定
説明このパラメーターをスケジューリングパラメーターと組み合わせて使用すると、毎日指定されたファイルからデータを読み取ることができます。
...
...
...
増分同期の設定
Data Integration のバッチ同期ノードでは、スケジューリングパラメーターを使用して、ソーステーブルとターゲットテーブルのデータパスと範囲を指定できます。設定方法は他のノードタイプと同じです。
実行時には、システムはノードで設定されたすべてのプレースホルダーパラメーターを、スケジューリングパラメーター式が表す実際の値に置き換えてから、データ同期を実行します。
MySQL のデータ同期を例に説明します。
データフィルタリング を設定しない場合、デフォルトですべてのデータがターゲットテーブルに同期されます。
データフィルタリング を設定した場合、フィルター条件を満たすデータのみがターゲットテーブルに同期されます。
ターゲットの MaxCompute テーブルのパーティション名は、スケジューリングパラメーターで指定します。$bizdate は業務日を表します。スケジュールされたタスクが実行されると、タスクに設定されたパーティション式は、スケジューリングパラメーターが表す業務日に置き換えられます。スケジューリングパラメーター式の詳細な設定手順については、「Data Integration におけるスケジューリングパラメーターの適用シナリオ」をご参照ください。バッチ同期タスクを例にとると、増分同期を実装するには、3 か所で bizdate パラメーターを設定する必要があります。ソースの [データフィルタリング] セクションで、STR_TO_DATE('${bizdate}','%Y%m%d') <= gmt_modify_time AND gmt_modify_time < DATE_ADD(STR_TO_DATE('${bizdate}','%Y%m%d'), interval 1 day) と入力して、業務日に変更されたデータをフィルタリングします。ターゲットの [パーティション情報] セクションで、pt=${bizdate} と入力して、対応する日付パーティションにデータを書き込み、[クリーンアップルール] を [書き込み前に既存のデータをクリーンアップする (Insert Overwrite)] に設定します。右側の [スケジュール設定] の [パラメーター] セクションで、bizdate=$bizdate と入力して、スケジューリングシステムが実行時に ${bizdate} を実際の業務日に自動的に置き換えるようにします。増分データ同期を設定する際:
時刻型列に基づく増分同期:スケジューリングパラメーターを使用して、時刻型データを動的に置き換えることができます。タスクのスケジューリング時に、スケジューリングパラメーターは業務日に基づいて自動的に特定の値に置き換えられます。スケジューリングパラメーターの詳細については、「スケジューリングパラメーターの設定」をご参照ください。
非時刻型列に基づく増分同期:代入ノードを使用して列をターゲットデータ型に変換し、その後 Data Integration に渡して同期できます。代入ノードの詳細については、「代入ノードの作成」をご参照ください。
注意事項
増分同期タスクを設定する際は、以下の点に注意してください:
書き込み前に既存のデータをクリーンアップする (Insert Overwrite) の安全性:複数の同期タスクが同じ MaxCompute テーブルの異なるパーティションに書き込む場合、Insert Overwrite 戦略は安全です。この戦略は、現在のタスクで指定されたパーティションデータのみをクリアし、テーブルの他のパーティションのデータには影響を与えないため、データの競合や誤削除を防ぎます。
パーティション範囲のバッチ上書きの制限:DataWorks では、バッチ上書きのパーティション設定で時間範囲 (
hh=00-23など) を指定することはサポートされていません。複数の時間分のデータを上書きするには、各時間に対して個別のタスクを設定してください。パーティションパラメーターは現在、単一の特定の値またはワイルドカード*のみをサポートしています。ワイルドカード構文:ソースに時間レベルのパーティションが含まれているが、ターゲットに日レベルのパーティションしかない場合、時間パーティションフィールドにワイルドカード
*を入力して、すべての時間データと一致させます。*は引用符なしで直接入力してください ("*"のようにしないでください)。そうしないと、構文エラーが発生します。
タイムスタンプベースの高頻度スケジュール済み増分同期
DataWorks では、バッチ同期タスクと定期スケジューリング (5 分ごとや 1 時間ごとなど) を組み合わせることで、タイムスタンプベースのスケジュール済み増分同期をサポートしています。このアプローチは、RDS MySQL から SelectDB や StarRocks などのターゲットへの T+1 またはニアリアルタイム同期シナリオに適しています。リアルタイム CDC タスクを必要とせず、SQL フィルタリングを通じて増分同期を実装するため、継続的に実行されるタスクのコストを回避できます。主な設定ポイント:
ソースの データフィルタリング の where 句で、タイムスタンプ列をフィルタリング条件として使用します。例:
gmt_modify_time >= '$[yyyymmddhhmiss-10/mi]' AND gmt_modify_time < '$[yyyymmddhhmiss]'定期スケジューリングパラメーター (
$[yyyymmddhhmiss]など) を設定して時間範囲を動的に計算し、各スケジューリング実行で指定された時間間隔内の増分データのみを同期するようにします。スケジューリングパラメーターの詳細な設定については、「Data Integration におけるスケジューリングパラメーターの適用シナリオ」をご参照ください。ターゲットの列マッピングで、パーティション列にマッピングする定数パラメーターを手動で追加し、動的パーティション書き込みを有効にします。
データベースレベルのバッチ同期タスクの増分設定
単一テーブルの同期タスクに加えて、データベースレベルのバッチ同期タスクを作成して、定期的な増分同期 を実装することもできます。タスクを作成する際に、増分同期を選択し、データベースレベルの同期タスクで 増分条件 を設定します。これにより、効率的な日レベルの増分パーティション同期 (例: create_time 列によるフィルタリング) が可能になります。
このアプローチは、複数のテーブルの同期を一元管理したいが、特定のテーブルに対してのみ増分処理が必要なシナリオに適しています。データベースレベルのバッチ同期タスクの完全な設定プロセスについては、「データベースレベルのバッチ同期タスクの設定」をご参照ください。
例
履歴データの同期:履歴の増分データをターゲットテーブルの対応する時刻パーティションに同期する必要がある場合は、オペレーションセンターのデータバックフィル機能を使用できます。データバックフィル機能の詳細については、「データバックフィル」をご参照ください。データ同期ノードの設定で、データソースとして MySQL を、ターゲットとして MaxCompute (ODPS) を選択し、テーブル名を
czdなどの値に設定します。データフィルタリング条件で、${bizdate}を使用して増分範囲を制御します (例:STR_TO_DATE('${bizdate}','%Y%m%d') <= gmt_modify_time)。パーティション情報をds=${bizdate}に設定し、クリーンアップルールを [書き込み前に既存のデータをクリーンアップする (Insert Overwrite)] に設定します。スケジュール設定のパラメーターセクションで、bizdate=$bizdateを定義します。このスケジューリングパラメーターは、バックフィル時に業務日に基づいて自動的に特定の日付値に置き換えられます。バックフィルを実行する際、複数の業務日範囲 (例: 2022-05-01 から 2022-05-31 および 2022-04-01 から 2022-04-30) を設定し、[スケジュール時刻が現在時刻より後のバックフィルインスタンスを即座に実行] を選択し、実行順序として [業務日の昇順] を選択できます。
よくある質問
MaxCompute Reader のパーティションフィルターでスケジューリングパラメーターを使用した後、パーティションが見つからないエラーが発生した場合はどうすればよいですか?
原因:設定された スケジューリングパラメーター が実行時に実際のパーティション値に正しく解決されていないか、解決された値がソーステーブルの実際のパーティションと一致していません。
解決策:Data Studio のパラメーター設定を確認し、以下の条件が満たされていることを確認してください:
設定されたパラメーター名が入力パラメーター名と一致していることを確認します。
渡されたパラメーター値が MaxCompute の実際のパーティションと完全に一致していることを確認します。
DataWorks のバッチ同期タスクは、デフォルトで完全同期と増分同期のどちらを実行しますか? パーティション列のないソーステーブルに対して増分同期を設定するにはどうすればよいですか?
デフォルトの動作
DataWorks Data Integration のバッチ同期タスクは、デフォルトで完全同期を実行します。つまり、毎回すべてのデータが同期されます。増分同期は、スケジューリングパラメーターと組み合わせた データフィルタリング 条件を設定した場合にのみ有効になります。
パーティション列のないソーステーブルの処理
ソースデータベーステーブル (RDS テーブルなど) に時刻またはパーティション列がない場合、where 句を使用して増分データを直接フィルタリングすることはできません。ソーステーブルに時刻列 (dt や gmt_modify_time など) を追加して、増分フィルタリングの基準とすることを推奨します。列の準備ができたら、本トピックの「増分同期の設定」セクションを参照して、増分同期ロジックを設定してください。
基本概念
パラメーターの定義
splitPk(分割キー):主キー列を指定します。DataWorks は、この列の値の範囲に基づいてデータを複数のチャンクに分割し、マルチスレッドの同時読み取りを可能にします。splitFactor(分割ファクター):分割の粒度を制御します。値が大きいほど、分割が細かくなり、読み取りスレッド数が増加します。
パフォーマンスへの影響と推奨事項
splitPk と splitFactor を有効にすると、ソースデータベースへの負荷が増加します。ソースへの負荷を軽減するには、以下を推奨します。
同時実行数を減らすか、データベースレベルのバッチ同期タスクの同時実行数を 1 に設定してください。
分割列 (
splitPk) にインデックスが設定されていることを確認すると、読み取り効率が向上します。
ビジネスに適した同期ソリューションを選択するにはどうすればよいですか?
同期ソリューションは、実行モード (スケジュール済みバッチ同期またはリアルタイム同期) とデータ範囲 (完全または増分) という 2 つの独立した次元によって決まります。これらの 2 つの次元は自由に組み合わせることができます。本トピックで説明するシナリオは、スケジュール済みバッチ同期 + 増分同期です。
実行モードの選択
デフォルトでは、スケジュール済みバッチ同期を推奨します。タスクが終了するとすぐにリソースが解放され、実際の実行時間に基づいて課金されます。以下の 2 つの場合にのみリアルタイム同期を選択してください:
ビジネスで秒レベルのレイテンシが必要な場合。
ソースからの物理的な DELETE 操作を同期する必要がある場合。
これらの場合、リアルタイム同期が必要です。なぜなら、増分バッチ同期はタイムスタンプなどの条件を使用してデータをフィルタリングするためです。ソースで物理的に削除されたレコードは、クエリ結果に表示されなくなるため、ターゲットでも削除されません。ここではリアルタイム同期について言及するに留め、その設定は本トピックでは扱いません。詳細については、「単一テーブルのリアルタイム同期機能」をご参照ください。
データ範囲の選択
完全同期は、データ量が少ないソーステーブル、または履歴データがその場で変更され、変更時刻列がないために変更範囲を特定できないシナリオに適しています。
増分同期では、ソーステーブルに信頼性の高い変更時刻列 (
gmt_modify_timeなど) または自動インクリメント主キーが必要です。大量のデータで日次変更率が小さく、日付パーティションに履歴スナップショットを保持する必要があるシナリオ(ジッパーテーブル、または緩やかに変化するディメンションテーブルとも呼ばれます)に適しています。ソーステーブルにそのような列がない場合は、本トピックの以下の FAQ をご参照ください: DataWorks のバッチ同期タスクは、デフォルトで完全同期と増分同期のどちらを実行しますか? パーティション列のないソーステーブルに対して増分同期を設定するにはどうすればよいですか?
初期ロードと大量データ
履歴データには「最初に完全、その後増分」の戦略を使用し、2 つの部分を異なるパーティションに書き込みます。
元のデータタイムスタンプで履歴パーティションをバックフィルするには、データバックフィル機能を使用してください。詳細については、「データバックフィルインスタンスの運用と保守」をご参照ください。
TB レベルのデータ量の場合、タスクの同時実行数を増やし、Data Integration リソースグループの CU を調整することで、同期効率を維持できます。詳細については、「サーバーレスリソースグループの作成と使用」をご参照ください。ただし、同時実行数を増やすと、ソースデータベースの負荷が増加します。詳細については、本トピックの以下の FAQ をご参照ください: 基本概念