public class DatahubClient extends Object
| 限定符和类型 | 类和说明 |
|---|---|
static class |
DatahubClient.ShardState
shard的状态
|
| 构造器和说明 |
|---|
DatahubClient(Odps odps,
String projectName,
String tableName,
String datahubEndpoint)
Datahub服务入口类
|
| 限定符和类型 | 方法和说明 |
|---|---|
String |
getProjectName() |
List<Long> |
getShardList() |
HashMap<Long,DatahubClient.ShardState> |
getShardStatus()
查询DatahubClinet对应的table拥有的shard在服务端的状态
|
TableSchema |
getStreamSchema() |
String |
getTableName() |
void |
loadShard(long shardNumber)
在ODPS hub服务上启用shard
|
DatahubReader |
openDatahubReader(long shardId)
创建DatahubReader读取指定shard
|
DatahubReader |
openDatahubReader(long shardId,
String packId)
创建DatahubReader读取指定shard
|
DatahubWriter |
openDatahubWriter()
创建DatahubWriter
|
DatahubWriter |
openDatahubWriter(long shardId)
创建DatahubWriter写入指定shard
|
PackReader |
openPackReader(long shardId) |
PackReader |
openPackReader(long shardId,
String packId) |
ReplicatorStatus |
QueryReplicatorStatus(long shardId)
在ODPS hub查询非分区表拷贝到离线集群的状态
|
ReplicatorStatus |
QueryReplicatorStatus(long shardId,
PartitionSpec partitionSpec)
在ODPS hub查询partiton对应的拷贝到离线集群的状态
|
void |
setEndpoint(String endpoint)
设置DatahubServer地址
没有设置DatahubServer地址的情况下, 自动选择
|
void |
waitForShardLoad()
同步等待 load shard 完成
默认超时时间为 120000ms
|
void |
waitForShardLoad(long timeout)
同步等待 load shard 完成
最大超时时间为 120000ms
|
public DatahubClient(Odps odps, String projectName, String tableName, String datahubEndpoint) throws OdpsException
odps - odps对象projectName - 对应project名称tableName - 对应table名称datahubEndpoint - datahub服务地址,公网用户使用 http://dh.odps.aliyun.com,ecs或内网用户请使用 http://dh-ext.odps.aliyun-inc.comOdpsExceptionpublic String getProjectName()
public String getTableName()
public void loadShard(long shardNumber)
throws OdpsException
shardNumber - 需要启用的shard数量OdpsExceptionpublic void waitForShardLoad()
throws OdpsException
OdpsExceptionpublic void waitForShardLoad(long timeout)
throws OdpsException
timeout - 超时时间,单位是毫秒
若该值超过 120000ms,将等待 120000msOdpsExceptionpublic HashMap<Long,DatahubClient.ShardState> getShardStatus() throws OdpsException, IOException
OdpsException, - IOExceptionOdpsExceptionIOExceptionpublic ReplicatorStatus QueryReplicatorStatus(long shardId, PartitionSpec partitionSpec) throws OdpsException
shardId - 需要查询的shardIdpartitionSpec - 查询的分区,分区表必选, 非分区表可以为nullOdpsExceptionpublic void setEndpoint(String endpoint) throws OdpsException
没有设置DatahubServer地址的情况下, 自动选择
endpoint - OdpsExceptionpublic ReplicatorStatus QueryReplicatorStatus(long shardId) throws OdpsException
shardId - 需要查询的shardIdOdpsExceptionpublic TableSchema getStreamSchema()
public DatahubWriter openDatahubWriter(long shardId) throws OdpsException, IOException
shardId - 需要写入数据的shardIdOdpsException, - IOExceptionOdpsExceptionIOExceptionpublic DatahubWriter openDatahubWriter() throws OdpsException, IOException
OdpsException, - IOExceptionOdpsExceptionIOExceptionpublic DatahubReader openDatahubReader(long shardId) throws OdpsException, IOException
shardId - 需要读取数据的shardIdOdpsException, - IOExceptionOdpsExceptionIOExceptionpublic DatahubReader openDatahubReader(long shardId, String packId) throws OdpsException, IOException
shardId - 需要读取数据的shardIdpackId - 指定读取的packIdOdpsException, - IOExceptionOdpsExceptionIOExceptionpublic PackReader openPackReader(long shardId) throws OdpsException, IOException
public PackReader openPackReader(long shardId, String packId) throws OdpsException, IOException
Copyright © 2015 Alibaba Cloud Computing. All rights reserved.