【问题标题】:Apache Thrift: Multitask single Server and ClientApache Thrift:多任务单服务器和客户端
【发布时间】:2017-07-23 07:09:57
【问题描述】:

我读过thisthis。但是,我的情况不同。我不需要服务器上的多路复用服务,也不需要到服务器的多个连接。

背景:
对于我的大数据项目,我需要计算给定大数据的核心集。 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


    【解决方案1】:

    从本质上看,这更像是一个架构问题,而不是 Thrift 问题。给定前提

    我不需要每个连接一个线程,我想要一个连接,并且服务器同时处理多个服务请求。可视化:

    我解决了客户端等待服务器完成计算的部分问题,通过在函数声明中添加oneway来表示客户端只发出请求,根本不等待任何响应。

    准确地描述了用例,你想要这个:

    +---------------------+
    | Client              |
    +---------+-----------+
              |
              |
    +---------v-----------+
    | Server              |
    +---------+-----------+
              |
              |
    +---------v-----------+          +---------------------+
    | Queue<WorkItems>    <----------+ Worker Thread Pool  |
    +---------------------+          +---------------------+
    

    服务器的唯一任务是获取请求并尽快将它们插入工作项队列。这些工作项由一个独立的工作线程池处理,否则它完全独立于服务器部分。唯一共享的部分是工作项队列,这当然需要正确同步的访问方法。

    关于 serevr 选择:如果服务器足够快地处理请求,即使是 TSimpleServer 也可以。

    【讨论】:

    • 好的,你的方法确实比我的好。服务器对工作进行排队,然后工作人员获取一个块并进行计算。请再问一个问题,如果我在客户端中对解析进行多任务处理(这很耗时,所以我想多任务处理),同一客户端中的不同线程可能同时拥有client.compressPoints(message)。这容易出问题吗?
    • 通常客户端不是线程安全的,底层物理连接句柄也不是。 IOW,每个线程一个客户端。如果客户端保持连接打开,那么您将需要足够数量的服务器端点,因此TSimpleSever 将不再适合这种情况。
    猜你喜欢
    • 2018-06-08
    • 2011-01-31
    • 2017-08-20
    • 2012-05-23
    • 2014-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-08-06
    相关资源
    最近更新 更多