什么是 SMA(简单移动平均线)? 一个简单的或算术移动平均线,通过将多个时间段的证券收盘价相加然后除以该总价计算得出时间段数。
例如在上面的例子中,收盘价是:37.14 (2008-02-29), 38.32 (2008-03-01), 38.00 (2008-03-02), 38.71 (2008-03-03), 38.37 (2008- 03-04), 36.60 (2008-03-05)。
所以 2008-03-02 的 3 天 SMA 是 (37.14 + 38.32 + 38.00) / 3 = 37.82
2008-02-29 没有 3 天 SMA(因为只有 1 天的数据:2008-02-29),2008-03-01 没有 3 天 SMA(只有 2 天的数据: 2008-02-29, 2008-03-01)。
以下是针对您的数据的 3 天 SMA 的解决方案(您可以轻松地将其更改为“n”天 SMA)。
映射器(m.py):
import sys
for line in sys.stdin:
val = line.strip()
vals = val.split('\t')
print "%s\t%s:%s" % (vals[0], vals[1], vals[2])
映射器逻辑:
它只是读取行中的制表符分隔值并输出“{key}\t{val1}:{val2}.
例如对于第一行(制表符分隔值):
2008-03-05 36.60 36.60
它输出:
2008-03-05 36.60:36.60
减速器(r.py):
import sys
lValueA = list()
lValueB = list()
smaInterval = 3
for line in sys.stdin:
(key, val) = line.strip().split('\t')
vals = val.split(':')
lValueA.append(float(vals[0]))
lValueB.append(float(vals[1]))
if len(lValueA) == smaInterval:
sumA = 0;
sumB = 0;
for a in lValueA:
sumA += a
for b in lValueB:
sumB += b
sumA = sumA / smaInterval;
sumB = sumB / smaInterval;
print "%s\t%.2f\t%.2f" % (key, sumA, sumB);
del lValueA[0]
del lValueB[0]
减速器逻辑:
- 它使用 2 个列表。一个用于库存 A,一个用于库存 B。
- 假设 SMA 区间为 3 (
smaInterval = 3)
- 当输入一行时,它会解析该行并将值 A 和值 B 附加到各自的列表中
- 当任何列表的大小达到 3(即 SMA 区间)时,它会计算移动平均值并输出(键,股票 A 的 SMA,股票 B 的 SMA)并从每个列表中删除第零个元素。
我为您的输入执行了此操作。
我执行了它,没有使用 Hadoop,如下所示(input.txt 包含您在问题中提到的输入,带有制表符分隔值):
cat input.txt | python m.py | sort | python r.py
我得到以下输出(我验证是正确的):
2008-03-02 37.82 37.82
2008-03-03 38.34 38.34
2008-03-04 38.36 38.36
2008-03-05 37.89 37.89
您应该能够使用 Hadoop 框架执行相同的操作:
hadoop jar hadoop-streaming-2.7.1.jar -input {Input directory in HDFS} -output {Output directory in HDFS} -mapper {Path to the m.py} -reducer {Path to the r.py}
注意:
这段代码可以优化,你可能根本不需要减速器。如果您的数据很小,您可以在映射器端读取所有值,对它们进行排序,然后计算 SMA。我刚刚编写了这段代码,以说明使用 Hadoop 流计算 SMA。