Tablestore は、テーブル間でデータを移行または同期するための複数の方法を提供します。トンネルサービス、DataWorks、DataX、またはコマンドラインインターフェイスを使用して、あるテーブルから別のテーブルにデータを同期できます。
前提条件
ソーステーブルとターゲットテーブルのインスタンス名、エンドポイント、およびリージョン ID を取得してください。
Alibaba Cloud アカウントまたは Tablestore の権限を持つ RAM ユーザーの AccessKey を作成してください。
SDK を使用したデータ同期
トンネルサービスを使用して、テーブル間でデータを同期できます。この方法は、同じリージョン内、異なるリージョン間、および異なるアカウント間でのデータ同期をサポートしています。トンネルサービスは、データの変更をキャプチャし、リアルタイムでターゲットテーブルに同期します。次の例は、Java SDK を使用してこの同期を実装する方法を示しています。
コードを実行する前に、ソーステーブルとターゲットテーブルの名前、インスタンス名、エンドポイントのプレースホルダーを実際の値に置き換えてください。また、AccessKey ID と AccessKey Secret を環境変数として設定する必要があります。
import com.alicloud.openservices.tablestore.*;
import com.alicloud.openservices.tablestore.core.auth.DefaultCredentials;
import com.alicloud.openservices.tablestore.core.auth.ServiceCredentials;
import com.alicloud.openservices.tablestore.model.*;
import com.alicloud.openservices.tablestore.model.tunnel.*;
import com.alicloud.openservices.tablestore.tunnel.worker.IChannelProcessor;
import com.alicloud.openservices.tablestore.tunnel.worker.ProcessRecordsInput;
import com.alicloud.openservices.tablestore.tunnel.worker.TunnelWorker;
import com.alicloud.openservices.tablestore.tunnel.worker.TunnelWorkerConfig;
import com.alicloud.openservices.tablestore.writer.RowWriteResult;
import com.alicloud.openservices.tablestore.writer.WriterConfig;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
public class TableSynchronization {
// ソーステーブルの設定:テーブル名、インスタンス名、エンドポイント、AccessKey ID、AccessKey Secret
final static String sourceTableName = "sourceTableName";
final static String sourceInstanceName = "sourceInstanceName";
final static String sourceEndpoint = "sourceEndpoint";
final static String sourceAccessKeyId = System.getenv("SOURCE_TABLESTORE_ACCESS_KEY_ID");
final static String sourceKeySecret = System.getenv("SOURCE_TABLESTORE_ACCESS_KEY_SECRET");
// ターゲットテーブルの設定:テーブル名、インスタンス名、エンドポイント、AccessKey ID、AccessKey Secret
final static String targetTableName = "targetTableName";
final static String targetInstanceName = "targetInstanceName";
final static String targetEndpoint = "targetEndpoint";
final static String targetAccessKeyId = System.getenv("TARGET_TABLESTORE_ACCESS_KEY_ID");
final static String targetKeySecret = System.getenv("TARGET_TABLESTORE_ACCESS_KEY_SECRET");
// トンネル名
static String tunnelName = "source_table_tunnel";
// TablestoreWriter:高同時実行データ書き込み用のツール。
static TableStoreWriter tableStoreWriter;
// 成功した行と失敗した行の統計
static AtomicLong succeedRows = new AtomicLong();
static AtomicLong failedRows = new AtomicLong();
public static void main(String[] args) {
// ターゲットテーブルを作成します。
createTargetTable();
System.out.println("Create target table Done.");
// TunnelClient を初期化します。
TunnelClient tunnelClient = new TunnelClient(sourceEndpoint, sourceAccessKeyId, sourceKeySecret, sourceInstanceName);
// データトンネルを作成します。
String tunnelId = createTunnel(tunnelClient);
System.out.println("Create tunnel Done.");
// TablestoreWriter を初期化します。
tableStoreWriter = createTablesStoreWriter();
// データトンネルを介してデータを同期します。
TunnelWorkerConfig config = new TunnelWorkerConfig(new SimpleProcessor());
TunnelWorker worker = new TunnelWorker(tunnelId, tunnelClient, config);
try {
System.out.println("Connect tunnel and working ...");
worker.connectAndWorking();
// トンネルのステータスを監視します。トンネルのステータスが完全データ同期から増分データ同期に変わると、データ同期が完了します。
while (true) {
if (tunnelClient.describeTunnel(new DescribeTunnelRequest(sourceTableName, tunnelName)).getTunnelInfo().getStage().equals(TunnelStage.ProcessStream)) {
break;
}
Thread.sleep(5000);
}
// 同期結果
System.out.println("Data Synchronization Completed.");
System.out.println("* Succeed Rows Count: " + succeedRows.get());
System.out.println("* Failed Rows Count: " + failedRows.get());
// トンネルを削除します。
tunnelClient.deleteTunnel(new DeleteTunnelRequest(sourceTableName, tunnelName));
// リソースをシャットダウンします。
worker.shutdown();
config.shutdown();
tunnelClient.shutdown();
tableStoreWriter.close();
}catch(Exception e){
e.printStackTrace();
worker.shutdown();
config.shutdown();
tunnelClient.shutdown();
tableStoreWriter.close();
}
}
private static void createTargetTable() throws ClientException {
// ソーステーブルの情報を照会します。
SyncClient sourceClient = new SyncClient(sourceEndpoint, sourceAccessKeyId, sourceKeySecret, sourceInstanceName);
DescribeTableResponse response = sourceClient.describeTable(new DescribeTableRequest(sourceTableName));
// ターゲットテーブルを作成します。
SyncClient targetClient = new SyncClient(targetEndpoint, targetAccessKeyId, targetKeySecret, targetInstanceName);
TableMeta tableMeta = new TableMeta(targetTableName);
response.getTableMeta().getPrimaryKeyList().forEach(
item -> tableMeta.addPrimaryKeyColumn(new PrimaryKeySchema(item.getName(), item.getType()))
);
TableOptions tableOptions = new TableOptions(-1, 1);
CreateTableRequest request = new CreateTableRequest(tableMeta, tableOptions);
targetClient.createTable(request);
// リソースをシャットダウンします。
sourceClient.shutdown();
targetClient.shutdown();
}
private static String createTunnel(TunnelClient client) {
// データトンネルを作成し、トンネル ID を返します。
CreateTunnelRequest request = new CreateTunnelRequest(sourceTableName, tunnelName, TunnelType.BaseAndStream);
CreateTunnelResponse response = client.createTunnel(request);
return response.getTunnelId();
}
private static class SimpleProcessor implements IChannelProcessor {
@Override
public void process(ProcessRecordsInput input) {
if(input.getRecords().isEmpty())
return;
System.out.print("* Begin consume " + input.getRecords().size() + " records ... ");
for (StreamRecord record : input.getRecords()) {
switch (record.getRecordType()) {
// 行データを書き込みます。
case PUT:
RowPutChange putChange = new RowPutChange(targetTableName, record.getPrimaryKey());
putChange.addColumns(getColumnsFromRecord(record));
tableStoreWriter.addRowChange(putChange);
break;
// 行データを更新します。
case UPDATE:
RowUpdateChange updateChange = new RowUpdateChange(targetTableName, record.getPrimaryKey());
for (RecordColumn column : record.getColumns()) {
switch (column.getColumnType()) {
// 属性列を追加します。
case PUT:
updateChange.put(column.getColumn().getName(), column.getColumn().getValue(), System.currentTimeMillis());
break;
// 属性列のバージョンを 1 つ削除します。
case DELETE_ONE_VERSION:
updateChange.deleteColumn(column.getColumn().getName(),
column.getColumn().getTimestamp());
break;
// 属性列を削除します。
case DELETE_ALL_VERSION:
updateChange.deleteColumns(column.getColumn().getName());
break;
default:
break;
}
}
tableStoreWriter.addRowChange(updateChange);
break;
// 行データを削除します。
case DELETE:
RowDeleteChange deleteChange = new RowDeleteChange(targetTableName, record.getPrimaryKey());
tableStoreWriter.addRowChange(deleteChange);
break;
}
}
// バッファからデータをフラッシュします。
tableStoreWriter.flush();
System.out.println("Done");
}
@Override
public void shutdown() {
}
}
public static List<Column> getColumnsFromRecord(StreamRecord record) {
List<Column> retColumns = new ArrayList<>();
for (RecordColumn recordColumn : record.getColumns()) {
// データバージョンを現在のタイムスタンプに置き換えて、最大バージョン偏差の超過を防ぎます。
Column column = new Column(recordColumn.getColumn().getName(), recordColumn.getColumn().getValue(), System.currentTimeMillis());
retColumns.add(column);
}
return retColumns;
}
private static TableStoreWriter createTablesStoreWriter() {
WriterConfig config = new WriterConfig();
// 行レベルのコールバックで、成功した行と失敗した行の数をカウントし、同期に失敗した行に関する情報を出力します。
TableStoreCallback<RowChange, RowWriteResult> resultCallback = new TableStoreCallback<RowChange, RowWriteResult>() {
@Override
public void onCompleted(RowChange rowChange, RowWriteResult rowWriteResult) {
succeedRows.incrementAndGet();
}
@Override
public void onFailed(RowChange rowChange, Exception exception) {
failedRows.incrementAndGet();
System.out.println("* Failed Rows: " + rowChange.getTableName() + " | " + rowChange.getPrimaryKey() + " | " + exception.getMessage());
}
};
ServiceCredentials credentials = new DefaultCredentials(targetAccessKeyId, targetKeySecret);
return new DefaultTableStoreWriter(targetEndpoint, credentials, targetInstanceName,
targetTableName, config, resultCallback);
}
}DataWorks を使用したデータ同期
DataWorks は、グラフィカルインターフェイスを通じて Tablestore テーブル間の同期タスクを設定できる、視覚的なデータ統合サービスを提供します。また、DataX などの他のツールを使用して、Tablestore テーブル間でデータを同期することもできます。
ステップ 1:準備
ターゲットテーブルを作成します。プライマリキー列のデータ型と順序を含むプライマリキー構造が、ソーステーブルと同じであることを確認してください。
DataWorks を有効化し、ソーステーブルまたはターゲットテーブルが配置されているリージョンに ワークスペースを作成します。
サーバーレスリソースグループを作成し、ワークスペースにバインドします。課金の詳細については、「サーバーレスリソースグループの課金」をご参照ください。
ソーステーブルとターゲットテーブルが異なるリージョンにある場合は、VPC ピアリング接続を作成して、リージョン間のネットワーク接続を確立する必要があります。
ステップ 2:Tablestore データソースの追加
ソーステーブルインスタンスとターゲットテーブルインスタンスの両方に Tablestore データソースを追加します。
DataWorks コンソールにログインします。 ターゲットリージョンに切り替えます。 左側のナビゲーションペインで、 を選択します。 ドロップダウンリストから目的のワークスペースを選択し、[データ統合に移動] をクリックします。
左側メニューで、[データソース] をクリックします。
[データソース] ページで、[データソースを追加] をクリックします。
[データソースを追加] ダイアログボックスで、データソースタイプとして [Tablestore] を検索して選択します。
[OTS データソースの追加] ダイアログボックスで、次の表の説明に従ってデータソースパラメーターを設定します。
パラメーター
説明
データソース名
データソースの名前。名前には、文字、数字、アンダースコア (_) を使用できますが、先頭に数字またはアンダースコア (_) は使用できません。
データソースの説明
データソースの簡単な説明。説明は 80 文字以内で入力してください。
リージョン
Tablestore インスタンスが配置されているリージョンを選択します。
[Table Store インスタンス名]
Tablestore インスタンスの名前。
エンドポイント
Tablestore インスタンスのエンドポイント。[VPC アドレス] を使用することを推奨します。
AccessKey ID
Alibaba Cloud アカウントまたは RAM ユーザーの AccessKey ID と AccessKey Secret。
AccessKey Secret
リソースグループの接続性をテストします。
リソースグループのデータソースへの接続性をテストする必要があります。リソースグループがデータソースに接続できない場合、同期タスクは実行できません。
接続設定 セクションで、リソースグループの [接続ステータス] 列にある [ネットワーク接続のテスト] をクリックします。
接続性テストに合格すると、[接続ステータス] は [接続済み] に変わります。[完了] をクリックします。データソースリストで新しいデータソースを表示できます。
接続性テストの結果が失敗の場合、接続診断ツール を使用してご自身で問題を解決できます。
ステップ 3:同期タスクの設定と実行
タスクノードの作成
[データ開発] ページに移動します。
DataWorks コンソールにログインします。
上部メニューで、リソースグループとリージョンを選択します。
左側のナビゲーションペインで、 を選択します。
目的のワークスペースを選択し、[データスタジオへ] をクリックします。
Data Studio コンソールの [データ開発] ページで、プロジェクトディレクトリ の横にある
アイコンにポインターを合わせ、 を選択します。[ノードを作成] ダイアログボックスで、[パス] を選択します。データソースと宛先の両方で Tablestore を選択します。名前を入力し、[OK] をクリックします。
同期タスクの設定
プロジェクトディレクトリ で、作成したバッチ同期ノードをクリックして開きます。同期タスクは、コードレス UI またはコードエディタで設定できます。
コードレス UI (デフォルト)
次の項目を設定します。
[データソース]:ソースと宛先のデータソースを選択します。
[実行中のリソース]: リソースグループを選択します。 リソースグループを選択すると、データソースの接続性が自動的にテストされます。
[データソース]:
[テーブル]: ドロップダウンリストから、ソーステーブルを選択します。
プライマリーキー範囲の開始: JSON 配列として指定される、データ読み取りの開始プライマリーキーです。
inf_minは負の無限大を表します。プライマリキーが
int型の列idとstring型の列nameで構成される場合、設定例は次のとおりです。特定のプライマリキー範囲
完全データ
[ { "type": "int", "value": "000" }, { "type": "string", "value": "aaa" } ][ { "type": "inf_min" }, { "type": "inf_min" } ]主キー範囲 (終了): 読み取るデータの終了主キーです。JSON 配列として指定し、
inf_maxは無限大を表します。プライマリキーが
int型の列idとstring型の列nameで構成される場合、設定例は次のとおりです。特定のプライマリキー範囲
完全データ
[ { "type": "int", "value": "999" }, { "type": "string", "value": "zzz" } ][ { "type": "inf_max" }, { "type": "inf_max" } ][チャンク構成情報]: JSON 配列で指定するカスタムの分割設定。通常、このパラメーターは設定しないことを推奨します (
[]に設定します)。Tablestore にホットパーティションが存在し、Tablestore Reader の自動分割ポリシーが有効にならない場合は、カスタム分割ルールを使用することを推奨します。分割ルールは、開始プライマリキーと終了プライマリキーの範囲内の分割点を指定します。すべてのプライマリキーではなく、分割キーのみを設定する必要があります。
[データ転送先]:
[テーブル]: ドロップダウンリストから、宛先テーブルを選択します。
[プライマリキー情報]: 宛先テーブルのプライマリキー情報です。値は JSON 配列形式である必要があります。
主キーに
int型の主キー列idが 1 つとstring型の主キー列nameが 1 つ含まれる場合、設定例は次のとおりです。[ { "name": "id", "type": "int" }, { "name": "name", "type": "string" } ][書き込みモード]: Tablestore にデータを書き込むモードです。次のモードがサポートされています:
PutRow:行データを書き込みます。ターゲット行が存在しない場合、新しい行が追加されます。ターゲット行が存在する場合、既存の行は上書きされます。
UpdateRow:行データを更新します。行が存在しない場合、新しい行が追加されます。行が存在する場合、リクエストに基づいて、この行の指定された列の値が追加、変更、または削除されます。
[宛先フィールドマッピング]: ソーステーブルから宛先テーブルへのフィールドマッピングを設定します。各行は 1 つのフィールドを表し、JSON 形式で指定されます。
[ソースフィールド]:値には、ソーステーブルの主キー情報が含まれている必要があります。
プライマリキーに
int型のプライマリキー列idとstring型のプライマリキー列nameがあり、属性列にint型のフィールドageがある場合、設定例は次のようになります。{"name":"id","type":"int"} {"name":"name","type":"string"} {"name":"age","type":"int"}[ターゲットフィールド]:値は、ターゲットテーブルのプライマリーキー情報を含める必要はありません。
プライマリキーが
int型のカラムidとstring型のカラムnameで構成され、属性カラムにint型のフィールドageが含まれる場合、設定例は次のとおりです。{"name":"age","type":"int"}
設定が完了したら、ページの上部にある保存をクリックします。
コードエディター
ページ上部のコードエディタをクリックします。表示されたページで、スクリプトを編集します。
次の例は、主キーがint型の列idとstring型の列nameで構成され、属性列にint型のフィールドageが含まれるテーブルの設定を示します。設定を構成する際には、サンプルスクリプト内のdatasourceとtableの名前を置き換えてください。
完全データ
{
"type": "job",
"version": "2.0",
"steps": [
{
"stepType": "ots",
"parameter": {
"datasource": "source_data",
"column": [
{
"name": "id",
"type": "int"
},
{
"name": "name",
"type": "string"
},
{
"name": "age",
"type": "int"
}
],
"range": {
"begin": [
{
"type": "inf_min"
},
{
"type": "inf_min"
}
],
"end": [
{
"type": "inf_max"
},
{
"type": "inf_max"
}
],
"split": []
},
"table": "source_table",
"newVersion": "true"
},
"name": "Reader",
"category": "reader"
},
{
"stepType": "ots",
"parameter": {
"datasource": "target_data",
"column": [
{
"name": "age",
"type": "int"
}
],
"writeMode": "UpdateRow",
"table": "target_table",
"newVersion": "true",
"primaryKey": [
{
"name": "id",
"type": "int"
},
{
"name": "name",
"type": "string"
}
]
},
"name": "Writer",
"category": "writer"
}
],
"setting": {
"errorLimit": {
"record": "0"
},
"speed": {
"concurrent": 2,
"throttle": false
}
},
"order": {
"hops": [
{
"from": "Reader",
"to": "Writer"
}
]
}
}特定のプライマリキー範囲
{
"type": "job",
"version": "2.0",
"steps": [
{
"stepType": "ots",
"parameter": {
"datasource": "source_data",
"column": [
{
"name": "id",
"type": "int"
},
{
"name": "name",
"type": "string"
},
{
"name": "age",
"type": "int"
}
],
"range": {
"begin": [
{
"type": "int",
"value": "000"
},
{
"type": "string",
"value": "aaa"
}
],
"end": [
{
"type": "int",
"value": "999"
},
{
"type": "string",
"value": "zzz"
}
],
"split": []
},
"table": "source_table",
"newVersion": "true"
},
"name": "Reader",
"category": "reader"
},
{
"stepType": "ots",
"parameter": {
"datasource": "target_data",
"column": [
{
"name": "age",
"type": "int"
}
],
"writeMode": "UpdateRow",
"table": "target_table",
"newVersion": "true",
"primaryKey": [
{
"name": "id",
"type": "int"
},
{
"name": "name",
"type": "string"
}
]
},
"name": "Writer",
"category": "writer"
}
],
"setting": {
"errorLimit": {
"record": "0"
},
"speed": {
"concurrent": 2,
"throttle": false
}
},
"order": {
"hops": [
{
"from": "Reader",
"to": "Writer"
}
]
}
}スクリプトを編集した後、ページの上部にある保存をクリックします。
同期タスクの実行
ページの上部にある [実行] をクリックして同期タスクを開始します。タスクを初めて実行するときは、[デバッグ設定] を確認する必要があります。
ステップ 4:同期結果の表示
タスクが実行された後、ログで実行ステータスを確認し、Tablestore コンソールで同期されたデータを確認します。
ページの下部でタスクの実行ステータスと結果を表示します。次のログ情報は、同期タスクが成功したことを示しています。
2025-11-18 11:16:23 INFO Shell run successfully! 2025-11-18 11:16:23 INFO Current task status: FINISH 2025-11-18 11:16:23 INFO Cost time is: 77.208sターゲットテーブルのデータを表示します。
Tablestore コンソールに移動します。上部メニューで、リソースグループとリージョンを選択します。
インスタンスエイリアスをクリックします。テーブルリスト ページで、対象のテーブルをクリックします。
データのクエリ をクリックして、ターゲットテーブルのデータを表示します。
CLI を使用したデータ同期
コマンドラインインターフェイスを使用するには、ソーステーブルからローカルの JSON ファイルに手動でデータをエクスポートし、そのファイルをターゲットテーブルにインポートする必要があります。この方法は、小規模なデータ移行にのみ適しています。
ステップ 1:準備
ターゲットテーブルを作成します。プライマリキー列の名前、データ型、順序を含むプライマリキー構造が、ソーステーブルと同じであることを確認してください。
コマンドラインインターフェイスをダウンロードしてください。
ステップ 2:ソーステーブルからのデータのエクスポート
コマンドラインインターフェイスを起動し、config コマンドを実行して、ソーステーブルが配置されているインスタンスのアクセス情報を設定します。詳細については、「ツールの起動とアクセス情報の設定」をご参照ください。
コマンドを実行する前に、endpoint、instance、id、および key を、ソーステーブルが配置されているインスタンスのエンドポイント、インスタンス名、AccessKey ID、および AccessKey Secret に置き換えてください。
config --endpoint https://myinstance.cn-hangzhou.ots.aliyuncs.com --instance myinstance --id NTSVL******************** --key 7NR2****************************************データをエクスポートします。
ソーステーブルを使用するには、
useコマンドを実行します。このトピックでは、source_tableを例として使用します。use --wc -t source_tableソーステーブルからローカルの JSON ファイルにデータをエクスポートします。詳細については、「データのエクスポート」をご参照ください。
scan -o /tmp/sourceData.json
ステップ 3:ターゲットテーブルへのデータのインポート
config コマンドを実行して、ターゲットテーブルが配置されているインスタンスのアクセス情報を設定します。
コマンドを実行する前に、endpoint、instance、id、および key を、ターゲットテーブルが配置されているインスタンスのエンドポイント、インスタンス名、AccessKey ID、および AccessKey Secret に置き換えてください。
config --endpoint https://myinstance.cn-hangzhou.ots.aliyuncs.com --instance myinstance --id NTSVL******************** --key 7NR2****************************************データをインポートします。
useコマンドを実行して、対象のテーブルを使用します。この例ではtarget_tableを使用します。use --wc -t target_tableローカルの JSON ファイルからターゲットテーブルにデータをインポートします。詳細については、「データのインポート」をご参照ください。
import -i /tmp/sourceData.json