【发布时间】:2017-07-23 07:09:57
【问题描述】:
我读过this 和this。但是,我的情况不同。我不需要服务器上的多路复用服务,也不需要到服务器的多个连接。
背景:
对于我的大数据项目,我需要计算给定大数据的核心集。
Coreset 是大数据的一个子集,它保留了大数据最重要的数学关系。
工作流程:
- 将海量数据分割成更小的块
- 客户端解析块并将其发送到服务器
- 服务器计算核心集并保存结果
我的问题:
整个事情作为单线程执行。
客户端解析一个块,然后等待服务器完成计算核心集,然后再解析另一个块,依此类推。
目标:
利用多处理。客户端同时解析多个块,并且对于每个compute coreset 请求,服务器任务一个线程来处理它。线程数有限的地方。像游泳池这样的东西。
我知道我需要使用与 TSimpleServer 不同的协议并转向 TThreadPoolServer 或 TThreadedServer。我只是无法决定选择哪一个,因为两者似乎都不适合我?
TThreadedServer 为每个客户端连接生成一个新线程,并且每个线程保持活动状态,直到客户端连接关闭。
在 TThreadedServer 中,每个客户端连接都有自己的专用服务器线程。客户端关闭连接后,服务器线程返回线程池以供重用。
我不需要每个连接一个线程,我想要一个连接,并且服务器同时处理多个服务请求。 可视化:
Client:
Thread1: parses(chunk1) --> Request compute coreset
Thread2: parses(chunk2) --> Request compute coreset
Thread3: parses(chunk3) --> Request compute coreset
Server: (Pool of 2 threads)
Thread1: Handle compute Coreset
Thread2: handle compute Coreset
.
.
Thread1 becomes available and handles another compute coreset
代码:
api.thrift:
struct CoresetPoint {
1: i32 row,
2: i32 dim,
}
struct CoresetAlgorithm {
1: string path,
}
struct CoresetWeightedPoint {
1: CoresetPoint point,
2: double weight,
}
struct CoresetPoints {
1: list<CoresetWeightedPoint> points,
}
service CoresetService {
void initialize(1:CoresetAlgorithm algorithm, 2:i32 coresetSize)
oneway void compressPoints(1:CoresetPoints message)
CoresetPoints getTotalCoreset()
}
服务器:(为了更好看,移除了实现)
class CoresetHandler:
def initialize(self, algorithm, coresetSize):
def _add(self, leveledSlice):
def compressPoints(self, message):
def getTotalCoreset(self):
if __name__ == '__main__':
logging.basicConfig()
handler = CoresetHandler()
processor = CoresetService.Processor(handler)
transport = TSocket.TServerSocket(port=9090)
tfactory = TTransport.TBufferedTransportFactory()
pfactory = TBinaryProtocol.TBinaryProtocolFactory()
server = TServer.TThreadedServer(processor, transport, tfactory, pfactory)
# You could do one of these for a multithreaded server
# server = TServer.TThreadedServer(processor, transport, tfactory, pfactory)
# server = TServer.TThreadPoolServer(processor, transport, tfactory, pfactory)
print 'Starting the server...'
server.serve()
print 'done.'
客户:
try:
# Make socket
transport = TSocket.TSocket('localhost', 9090)
# Buffering is critical. Raw sockets are very slow
transport = TTransport.TBufferedTransport(transport)
# Wrap in a protocol
protocol = TBinaryProtocol.TBinaryProtocol(transport)
# Create a client to use the protocol encoder
client = CoresetService.Client(protocol)
# Connect!
transport.open()
// Here data is sliced, and in a loop I move on all files
Saved in the directory I specified, then they are parsed and
client.compressPoints(data) is invoked.
SliceFile(...)
p = CoresetAlgorithm(...)
client.initialize(p, 200)
for filename in os.listdir('/home/tony/DanLab/slicedFiles'):
if filename.endswith(".txt"):
data = _parse(filename)
client.compressPoints(data)
compressedData = client.getTotalCoreset()
# Close!
transport.close()
except Thrift.TException, tx:
print '%s' % (tx.message)
问题:
可以在 Thrift 中使用吗?我应该使用什么协议?
我通过在函数声明中添加oneway解决了客户端等待服务器完成计算的部分问题
to 表示客户端只发出请求,根本不等待任何响应。
【问题讨论】:
标签: python multithreading communication thrift