【发布时间】:2017-03-30 01:26:50
【问题描述】:
我有大量的网络数据,我一直使用这些数据进行 Spark 和 MLLib 的聚类练习。我已将数据标准化为一组向量,这些向量表示一天中的时间、方向(进出网络)、发送的字节数、接收的字节数以及每个连接的持续时间。一共有七个维度。
使用 KMeans,很容易用这些数据构建模型。使用这个模型,每个输入向量都被“分类”,距离被计算到最近的质心。最后,RDD(现在用距离标记)按距离排序,提取出最极端的值。
我的数据中的一个输入列是连接 uuid(唯一的字母数字标识符)。我很想通过模型携带这些数据(让每个输入向量都被唯一标记),但是当列无法转换为浮点数时会触发异常。
这里的问题是:“我如何最有效地将异常值与原始输入数据联系起来?”输入数据经过高度标准化,与原始输入不同。此外,源 IP 地址和目标 IP 地址已丢失。我在 KMeans 中没有看到任何接口来告诉它在构建模型时要考虑(或者相反地,忽略)哪些列。
我的代码如下所示:
def get_distance(clusters):
def _distance_map(record):
cluster = clusters.predict(record)
centroid = clusters.clusterCenters[cluster]
dist = np.linalg.norm(np.array(record) - np.array(centroid))
return (dist, record)
return _distance_map
def parseMap(row):
# parses rows of data out of the input strings
def conMap(row):
# normalizes the values to be used in building the model
rdd = sc.textFile('/data2/network/201610').filter(lambda r: r[0] != '#')
tcp = rdd.map(parseMap).filter(lambda r: r['proto'] == 'tcp')
cons = tcp.map(conMap) # this normalizes connection data
model = KMeans.train(cons, (24 * 7), maxIterations=25,
runs=1, initializationMode = "random")
data_distance = cons.map(get_distance(model)).sortByKey(ascending=False)
print(data_distance.take(10))
【问题讨论】:
标签: python apache-spark apache-spark-mllib