【问题标题】:Generator for n-records of a real time stream用于实时流的 n 条记录的生成器
【发布时间】:2018-07-10 03:29:08
【问题描述】:

我订阅了一个实时流,它以较慢的速度(每 1-5 秒 0.5 KB)发布一个小的 JSON 记录。发布者提供了一个暴露这些记录的 python 客户端。我将这些记录写入内存中的列表。客户端只是一个 python 包装器,用于在数据集的 HTTPS 端点上执行 curl 命令。数据集由过滤器和字段定义。我可以让客户离开几天,然后在午夜停止它,将多天的数据作为一批处理。

我想通过将流视为生成器来编写每个 n 条记录,而不是上面描述的多天批处理。客户端代码如下。我刚刚添加了 append() 行来创建一个名为“records”的列表(在内存中)以便稍后播放:

records=[]
data_set = api.get_dataset(dataset_id='abc')
for record in data_set.request_realtime(): 
    records.append(record)

正如预期的那样,在 Jupyter Notebook 中给了我 [*];并继续运行。

然后,我从内存中的列表中创建了一个生成器,用于提取一条记录(n=1 用于初始测试):

def Generator():
    count = 1
    while count < 2:
            for r in records:
                yield r.data
            count +=1

但是我的生成器定义也给了我 [*] 并继续计算;我理解这是因为该列表仍在内存中写入。但我认为我的生成器将能够锁定我的列表状态并产生前 n 条记录。但它没有。在这种情况下,我如何编码我的生成器?如果在这个用例中生成器不是一个好的选择,请提出建议。

为了给你完整的画面,如果我的代码正常工作,那么,我会实例化它,打印它,并按预期接收一个对象,如下所示:

>>>my_generator = Generator()
>>>print(my_generator)
<generator object Gen at 0x0000000009910510>

然后,我会将其写入 csv 文件,如下所示:

with open('myfile.txt', 'w') as f:
    cf = csv.DictWriter(f, column_headers, extrasaction='ignore') 
    cf.writeheader()
    cf.writerows(i.data for i in my_generator)

注意:我知道有很多工具可以做到这一点,例如卡夫卡;但我处于初始 PoC 阶段。请使用 Python 2x。一旦我的代码工作起来,我计划堆叠生成器来设置我的下一个 n 记录提取,这样我就不会在两者之间丢失数据。任何关于堆叠的指导也将不胜感激。

【问题讨论】:

    标签: python concurrency stream generator


    【解决方案1】:

    这不是并发的工作方式。除非您使用了一些您没有告诉我们的魔法,否则当您的第一个代码返回 * 时,您将无法运行更多代码。将生成器放在另一个单元格中只是将其添加到队列中以在第一个代码完成时运行 - 因为第一个代码永远不会完成,第二个代码甚至永远不会开始运行!

    我建议研究一些异步网络库,例如 asynciotwistedtrio。它们允许您使函数协作,因此当其中一个在等待数据时,另一个可以运行,而不是阻塞。您可能还必须将 api.get_dataset 代码重写为异步代码。

    【讨论】:

      猜你喜欢
      • 2011-04-24
      • 1970-01-01
      • 1970-01-01
      • 2017-01-02
      • 1970-01-01
      • 2010-09-23
      • 2013-04-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多