1.介绍
数据复制服务(Data Replication Service,简称DRS)是一种易用、稳定、高效、用于数据库实时迁移和数据库实时同步的云服务。 数据复制服务围绕云数据库,降低了数据库之间数据流通的复杂性,有效地帮助您减少数据传输的成本。 您可通过数据复制服务快速解决多场景下,数据库之间的数据流通问题,以满足数据传输业务需求。 数据复制服务提供了实时迁移、备份迁移、实时同步、数据订阅和实时灾备等多种功能。该示例展示了如何通过Java版SDK暂停DRS任务、续传DRS任务。
2.前置条件
- 1、获取华为云开发工具包(SDK),您也可以查看安装Python SDK。
- 2、要使用华为云 Python SDK,您需要拥有华为云账号以及该账号对应的 Access Key(AK)和 Secret Access Key(SK)。具体请参见:访问秘钥
- 3、已具备开发环境 ,支持Python3及其以上版本。
- 4、限速只对全量迁移阶段生效,增量迁移阶段不生效。
3.SDK获取和安装
您可以通过加入相应的依赖项方式获取和安装SDK,具体的SDK版本号请参见 SDK开发中心 。
pip install huaweicloudsdkdrs
4.关键代码片段
import json
import os
from huaweicloudsdkcore.auth.credentials import BasicCredentials
from huaweicloudsdkcore.http.http_config import HttpConfig
from huaweicloudsdkdrs.v3.drs_client import DrsClient
from huaweicloudsdkdrs.v3.model.batch_list_job_details_request import BatchListJobDetailsRequest
from huaweicloudsdkdrs.v3.model.batch_list_job_status_request import BatchListJobStatusRequest
from huaweicloudsdkdrs.v3.model.batch_list_progresses_request import BatchListProgressesRequest
from huaweicloudsdkdrs.v3.model.batch_pause_job_req import BatchPauseJobReq
from huaweicloudsdkdrs.v3.model.batch_query_job_req_page import BatchQueryJobReqPage
from huaweicloudsdkdrs.v3.model.batch_query_progress_req import BatchQueryProgressReq
from huaweicloudsdkdrs.v3.model.batch_restore_task_request import BatchRestoreTaskRequest
from huaweicloudsdkdrs.v3.model.batch_retry_req import BatchRetryReq
from huaweicloudsdkdrs.v3.model.batch_stop_jobs_request import BatchStopJobsRequest
from huaweicloudsdkdrs.v3.model.page_req import PageReq
from huaweicloudsdkdrs.v3.model.pause_info import PauseInfo
from huaweicloudsdkdrs.v3.model.query_jobs_req import QueryJobsReq
from huaweicloudsdkdrs.v3.model.retry_info import RetryInfo
from huaweicloudsdkdrs.v3.model.show_job_list_request import ShowJobListRequest
from huaweicloudsdkdrs.v3.region.drs_region import DrsRegion
import time
class PauseAndRetryJobDemo:
def __init__(self):
pass
@staticmethod
def main(args):
# 创建认证
# 认证用的ak和sk直接写到代码中有很大的安全风险,建议在配置文件或者环境变量中密文存放,使用时解密,确保安全;
# 本示例以ak和sk保存在环境变量中来实现身份验证为例,运行本示例前请先在本地环境中设置环境变量HUAWEICLOUD_SDK_AK和HUAWEICLOUD_SDK_SK。
ak = os.environ["HUAWEICLOUD_SDK_AK"]
sk = os.environ["HUAWEICLOUD_SDK_SK"]
auth = BasicCredentials(
ak=ak,
sk=sk,
)
# 配置客户端属性
config = HttpConfig.get_default_config()
config.ignore_ssl_verification = True
config.proxy_protocol = "http"
# 创建DrsClient实例
client = DrsClient.new_builder() \
.with_credentials(credentials=auth) \
.with_http_config(config=config) \
.with_region(region=DrsRegion.CN_NORTH_4) \
.build()
# 查询租户任务列表
job_infos = PauseAndRetryJobDemo.__show_job_list(client)
if (job_infos is None or len(job_infos) == 0):
return
# 批量查询任务详情
query_job_resp_list = PauseAndRetryJobDemo.__batch_list_job_details(client)
if (query_job_resp_list is None or len(query_job_resp_list) == 0):
return
# 暂停任务
if (PauseAndRetryJobDemo.__pause_job_fail(client)):
return
status = PauseAndRetryJobDemo.__get_status(client)
# 任务状态为暂停中时,后续可根据需要续传
PauseAndRetryJobDemo.__retry_job(client, "")
# 批量查询任务进度
PauseAndRetryJobDemo.__batch_list_progresses(client)
@staticmethod
def __show_job_list(client):
"""
查询租户任务列表
@param client
@return
"""
request = ShowJobListRequest()
query_jobs_req = QueryJobsReq()
# 名称或ID,选填
query_jobs_req.name = "<YOUR JOB NAME OR JOB id>"
# 引擎类型,选填
query_jobs_req.engine_type = "mysql"
# 任务类型,选填
query_jobs_req.db_use_type = "sync"
# 企业项目,选填
query_jobs_req.enterprise_project_id = "<YOUR JOB ENTERPRISE PROJECT ID>"
# 网络类型,选填
query_jobs_req.net_type = "eip"
# 服务名称,选填
query_jobs_req.service_name = "<YOUR SERVICE NAME>"
# 状态,选填
query_jobs_req.status = "CONFIGURATION"
# 标签,选填
tags = {}
query_jobs_req.tags = tags
# 每页记录数,选填, 默认10
query_jobs_req.per_page = 10
# 第几页,选填,默认1
query_jobs_req.cur_page = 1
request.body = query_jobs_req
request.x_language = "en-us"
show_job_list_response = client.show_job_list(request)
print(show_job_list_response.jobs)
return show_job_list_response.jobs
@staticmethod
def __batch_list_job_details(client):
"""
批量查询任务详情
@param client
@return
"""
batch_list_job_details_request = BatchListJobDetailsRequest()
batch_query_job_req_page = BatchQueryJobReqPage()
# 批量查询任务信息任务ID请求列表
jobs = []
jobs.append("<YOUR JOB ID>")
batch_query_job_req_page.jobs = jobs
# 分页请求体
page_req = PageReq()
# 当前页,选填,默认1
page_req.cur_page = 1
# 每页显示项数,选填,默认5
page_req.per_page = 5
batch_query_job_req_page.page_req = page_req
batch_list_job_details_request.body = batch_query_job_req_page
batch_list_job_details_response = client.batch_list_job_details(batch_list_job_details_request)
print(batch_list_job_details_response.results)
return batch_list_job_details_response.results
@staticmethod
def __pause_job_fail(client):
"""
暂停任务
@param client
@return
"""
batch_stop_jobs_request = BatchStopJobsRequest()
batch_pause_job_req = BatchPauseJobReq()
jobs = []
pause_info = PauseInfo()
pause_info.job_id = "<YOUR JOB ID>"
pause_info.pause_mode = "target"
jobs.append(pause_info)
batch_pause_job_req.jobs = jobs
batch_stop_jobs_request.body = batch_pause_job_req
batch_stop_jobs_response = client.batch_stop_jobs(batch_stop_jobs_request)
print(batch_stop_jobs_response)
if (batch_stop_jobs_response.http_status_code != 202):
print(batch_stop_jobs_response)
return True
return False
@staticmethod
def __retry_job(client, status):
"""
续传任务
@param client
"""
current_time = time.time()
while ("PAUSING" != status):
aaa = int(round(time.time()))
bbb = int(round(current_time))
if (int(round(time.time())) - int(round(current_time)) < 20):
continue
current_time = time.time()
status = PauseAndRetryJobDemo.__get_status(client)
batch_restore_task_request = BatchRestoreTaskRequest()
batch_retry_req = BatchRetryReq()
retry_infos = []
retry_info = RetryInfo()
retry_info.job_id = "<YOUR JOB ID>"
retry_infos.append(retry_info)
batch_retry_req.jobs = retry_infos
batch_restore_task_request.body = batch_retry_req
batch_restore_task_response = client.batch_restore_task(batch_restore_task_request)
print(batch_restore_task_response)
@staticmethod
def __get_status(client):
"""
获取任务状态
@param client
@return
"""
batch_list_progresses_request = BatchListJobStatusRequest()
batch_query_job_req_page = BatchQueryJobReqPage()
jobs = []
jobs.append("<YOUR JOB ID>")
batch_query_job_req_page.jobs = jobs
batch_list_progresses_request.body = batch_query_job_req_page
batch_list_job_status_response = client.batch_list_job_status(batch_list_progresses_request)
return batch_list_job_status_response.results[0].status
@staticmethod
def __batch_list_progresses(client):
"""
批量查询任务进度
@param client
@return
"""
batch_list_progresses_request = BatchListProgressesRequest()
batch_query_job_req_page = BatchQueryProgressReq()
job_ids = []
job_ids.append("<YOUR JOB ID>")
batch_query_job_req_page.jobs = job_ids
batch_list_progresses_request.body = batch_query_job_req_page
batch_list_progresses_response = client.batch_list_progresses(batch_list_progresses_request)
print(batch_list_progresses_response)
return batch_list_progresses_response.results
if __name__ == "__main__":
PauseAndRetryJobDemo().main(any)
5.参考链接
更多详细信息请参考:
6.修订记录
| 发布日期 | 文档版本 | 修订说明 |
|---|---|---|
| 2023-08-30 | 1.0 | 文档首次发布 |