氚云全量数据同步方案指南
概述与背景
当外部系统需要从氚云同步一张数据表的全部数据时,通常会调用 OpenAPI 的 LoadBizObjects 接口进行分页查询。当数据量较大(例如 30 万条)时,使用传统的页码分页(如 LIMIT 200000, 500)会导致数据库深度翻页(深分页),查询性能急剧下降,并消耗大量数据库资源。
本文档提供两种分页策略,旨在帮助开发者根据数据量选择最优方案,并提供了完整的、可直接参考的代码实现。
分页策略选型
| 策略名称 | 适用场景 | 核心原理 | 性能特点 |
|---|---|---|---|
| 深分页 | 数据总量 ≤ 1 万条 | 使用 FromRowNum/ToRowNum 进行传统页码偏移。 |
实现简单,但在数据量较大时,越往后查询越慢。 |
| 游标分页 | 数据总量 > 1 万条(推荐) | 利用索引字段(CreatedTime + ObjectId)记录上一页最后一条记录的位置,作为下一页查询的起点。 |
查询性能稳定,不受数据总量影响,是处理大数据量的标准做法。 |
核心概念:这里的“游标”并非数据库服务器端的游标,而是指客户端记录并传递的、用于标记上一页查询边界的一对字段值(即 lastCreatedTime 和 lastObjectId)。
方案一:深分页(适用于小数据量)
此方案通过控制页码(pageIndex)和每页大小(pageSize)来遍历数据。
关键代码说明
LoadPage方法:核心是构造Filter对象。FromRowNum=pageIndex * pageSize(起始行号)。ToRowNum=FromRowNum + pageSize(结束行号)。ReturnItems = new string[] { }表示返回所有字段(氚云 API 约定)。
- 循环终止条件:当 API 返回的
BizObjectArray数量小于PageSize时,表示已无更多数据,循环结束。
完整代码实现
public class LoadBizObjects_Deep_Page
{
private const string SchemaCode = "D00018smb9a4dbd53baf4950a6dafd6a92dff4b2";
private const int PageSize = 500; // 建议设为 500,平衡性能与请求次数
public List<Dictionary<string, object>> LoadBizObjects()
{
int pageIndex = 0;
var allBizObjects = new List<Dictionary<string, object>>();
while (true)
{
// 1. 请求当前页数据
string responseJson = LoadPage(pageIndex, PageSize);
var pageResult = JsonConvert.DeserializeObject<H3yunApiResponse<LoadBizObjectsReponseReturnData<Dictionary<string, object>>>>(responseJson);
// 2. 检查API调用是否成功
if (pageResult == null || !pageResult.Successful)
{
// 实际项目中应记录日志并抛出异常或返回已有数据
throw new Exception($"API调用失败: {pageResult?.ErrorMessage}");
}
var bizObjects = pageResult.ReturnData?.BizObjectArray;
if (bizObjects == null || bizObjects.Count == 0)
{
break; // 无数据,退出循环
}
allBizObjects.AddRange(bizObjects);
// 3. 判断是否为最后一页
if (bizObjects.Count < PageSize)
{
break;
}
pageIndex++; // 准备请求下一页
}
return allBizObjects;
}
private string LoadPage(int pageIndex, int pageSize)
{
var filter = new Filter
{
Matcher = new AndMatcher(),
FromRowNum = pageIndex * pageSize,
ToRowNum = (pageIndex + 1) * pageSize, // 更清晰的写法
ReturnItems = new string[] { } // 空数组表示返回所有字段
};
string filterJson = JsonConvert.SerializeObject(filter);
var parameters = new Dictionary<string, object>
{
["ActionName"] = "LoadBizObjects",
["SchemaCode"] = SchemaCode,
["Filter"] = filterJson
};
return HttpClientHelper.Post("https://www.h3yun.com/OpenApi/Invoke", parameters);
}
}
方案二:游标分页(推荐用于大数据量)
此方案通过排序字段(CreatedTime 升序,ObjectId 升序)确保数据顺序固定,然后使用上一页最后一条记录的这两个字段值作为下一页的查询起点。
核心逻辑解析(重点)
LoadNextPage 方法中构造的 Filter 是实现游标分页的关键。其查询条件 (CreatedTime, ObjectId) > (lastCreatedTime, lastObjectId) 在 SQL 中大致等价于:
SELECT * FROM TABLENAME
WHERE (CreatedTime > 'lastCreatedTime')
OR (CreatedTime = 'lastCreatedTime' AND ObjectId > 'lastObjectId')
ORDER BY CreatedTime ASC, ObjectId ASC
为什么用两个字段? CreatedTime 可能不是唯一的(同一秒内可能创建多条记录),因此需要组合 ObjectId(全局唯一)来保证排序的绝对稳定,避免数据遗漏或重复。
完整代码实现
public class LoadBizObjects_Cursor_Page
{
private const string SchemaCode = "D00018smb9a4dbd53baf4950a6dafd6a92dff4b2";
private const int PageSize = 500;
public List<Dictionary<string, object>> LoadBizObjects()
{
var allBizObjects = new List<Dictionary<string, object>>();
string firstPageJson = LoadFirstPage(PageSize);
var pageResult = JsonConvert.DeserializeObject<H3yunApiResponse<LoadBizObjectsReponseReturnData<Dictionary<string, object>>>>(firstPageJson);
if (pageResult == null || !pageResult.Successful)
{
throw new Exception($"首次加载失败: {pageResult?.ErrorMessage}");
}
var currentPageData = pageResult.ReturnData?.BizObjectArray;
if (currentPageData == null || currentPageData.Count == 0)
{
return allBizObjects; // 无数据
}
allBizObjects.AddRange(currentPageData);
// 循环获取后续页面
while (currentPageData.Count == PageSize) // 如果当前页数量等于页大小,说明可能还有下一页
{
// 获取当前页最后一条记录的“游标”值
var lastItem = currentPageData[currentPageData.Count - 1];
string lastObjectId = lastItem["ObjectId"]?.ToString();
string lastCreatedTime = lastItem["CreatedTime"]?.ToString();
// 请求下一页
string nextPageJson = LoadNextPage(lastObjectId, lastCreatedTime, PageSize);
pageResult = JsonConvert.DeserializeObject<H3yunApiResponse<LoadBizObjectsReponseReturnData<Dictionary<string, object>>>>(nextPageJson);
if (pageResult == null || !pageResult.Successful)
{
throw new Exception($"加载下一页失败: {pageResult?.ErrorMessage}");
}
currentPageData = pageResult.ReturnData?.BizObjectArray;
if (currentPageData == null || currentPageData.Count == 0)
{
break; // 无更多数据
}
allBizObjects.AddRange(currentPageData);
}
return allBizObjects;
}
private string LoadFirstPage(int pageSize)
{
Filter filter = new Filter
{
Matcher = new AndMatcher(),
FromRowNum = 0,
ToRowNum = pageSize,
ReturnItems = new string[] { },
// 关键:必须指定排序,且排序字段必须与游标条件字段一致
SortByCollection = JsonConvert.SerializeObject(new[]
{
new SortBy { ItemName = "CreatedTime", Direction = 0 }, // 0=升序
new SortBy { ItemName = "ObjectId", Direction = 0 }
})
};
string filterJson = JsonConvert.SerializeObject(filter);
Dictionary<string, object> parameters = new Dictionary<string, object>
{
["ActionName"] = "LoadBizObjects",
["SchemaCode"] = SchemaCode,
["Filter"] = filterJson
};
return HttpClientHelper.Post(OpenAPIConfig.OpenAPI_Url, parameters);
}
private string LoadNextPage(string lastObjectId, string lastCreatedTime, int pageSize)
{
// 构造条件:(CreatedTime > lastCreatedTime) OR (CreatedTime = lastCreatedTime AND ObjectId > lastObjectId)
OrMatcher orMatcher = new OrMatcher();
// 条件1:CreatedTime > lastCreatedTime
orMatcher.Matchers.Add(new FieldMatcher
{
Name = "CreatedTime",
Operator = 0, // 0 表示 "大于"
Value = lastCreatedTime
});
// 条件2:CreatedTime = lastCreatedTime AND ObjectId > lastObjectId
AndMatcher andMatcher = new AndMatcher();
andMatcher.Matchers.Add(new FieldMatcher
{
Name = "CreatedTime",
Operator = 2, // 2 表示 "等于"
Value = lastCreatedTime
});
andMatcher.Matchers.Add(new FieldMatcher
{
Name = "ObjectId",
Operator = 0, // 注意:这里的 0 表示 "大于"
Value = lastObjectId
});
orMatcher.Matchers.Add(andMatcher);
var filter = new Filter
{
Matcher = orMatcher,
FromRowNum = 0,
ToRowNum = pageSize,
ReturnItems = new string[] { },
SortByCollection = JsonConvert.SerializeObject(new[]
{
new SortBy { ItemName = "CreatedTime", Direction = 0 },
new SortBy { ItemName = "ObjectId", Direction = 0 }
})
};
string filterJson = JsonConvert.SerializeObject(filter);
var parameters = new Dictionary<string, object>
{
["ActionName"] = "LoadBizObjects",
["SchemaCode"] = SchemaCode,
["Filter"] = filterJson
};
return HttpClientHelper.Post("https://www.h3yun.com/OpenApi/Invoke", parameters);
}
}
通用组件与配置
API 响应结构
public class H3yunApiResponse<T> where T : class
{
public H3yunApiResponse()
{
Logined = false;
}
public bool Successful
{
get
{
return string.IsNullOrEmpty(ErrorMessage);
}
}
public string ErrorMessage { get; set; }
public bool Logined { get; set; }
public T ReturnData { get; set; }
}
public class LoadBizObjectsReponseReturnData<T> where T : class
{
public List<T> BizObjectArray { get; set; }
}
HTTP 帮助类与配置
HttpClientHelper:代码已提供,建议根据项目实际情况,将其中的HttpWebRequest替换为IHttpClientFactory以提升性能和资源管理。OpenAPIConfig:此类需开发者自行创建,用于存储氚云环境的访问地址、EngineCode和EngineSecret。
/// <summary>
/// HTTP请求帮助类
/// </summary>
public class HttpClientHelper
{
/// <summary>
/// 执行POST请求并返回JSON响应
/// </summary>
public static string Post(string url, Dictionary<string, object> parameters, string engineCode, string engineSecret)
{
HttpWebRequest request = (HttpWebRequest)WebRequest.Create(url);
request.Method = "POST";
request.ContentType = "application/json";
// 添加身份认证参数
request.Headers.Add("EngineCode", engineCode);
request.Headers.Add("EngineSecret", engineSecret);
// 序列化请求参数
string jsonData = JsonConvert.SerializeObject(parameters);
byte[] bytes = Encoding.UTF8.GetBytes(jsonData);
request.ContentLength = bytes.Length;
// 写入请求体
using (Stream writer = request.GetRequestStream())
{
writer.Write(bytes, 0, bytes.Length);
}
// 读取响应
StringBuilder returnStrBuilder = new StringBuilder();
using (HttpWebResponse response = (HttpWebResponse)request.GetResponse())
using (Stream s = response.GetResponseStream())
using (StreamReader reader = new StreamReader(s, Encoding.UTF8))
{
string lineStr;
while ((lineStr = reader.ReadLine()) != null)
{
returnStrBuilder.Append(lineStr);
returnStrBuilder.Append("\r\n");
}
}
return returnStrBuilder.ToString();
}
/// <summary>
/// 执行POST请求(使用默认配置)
/// </summary>
public static string Post(string url, Dictionary<string, object> parameters)
{
return Post(url, parameters, {EngineCode}, {EngineSecret});
}
}
最佳实践与注意事项
- 务必检查 API 响应状态:每次调用
LoadBizObjects后,都应检查pageResult.Successful,并对错误进行适当处理(如重试、记录日志、抛出异常)。 - 选择合适的分页大小:
PageSize建议设为 500,这是在请求次数和单次查询负载之间的一个良好平衡点。 - 确保排序稳定性:在使用游标分页时,
SortBy必须包含CreatedTime和ObjectId两个字段,并且顺序必须与WHERE条件中的逻辑严格对应。 - 考虑内存占用:本示例将所有数据加载到内存(
List<T>)中返回。如果数据量极大(如超过 50 万),建议改为流式处理或分批次写入目标系统,避免内存溢出。