【发布时间】: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