###spark版### Spark Graphx 进行团伙的识别(community detection)

2024-05-07 14:58

本文主要是介绍###spark版### Spark Graphx 进行团伙的识别(community detection),希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

最近在使用Spark Graphx,拿Graphx做了点实验。对大规模图常见的分析方法有连通图挖掘,团伙挖掘等。在金融科技领域,尤其风控领域,会有各种重要的关联网络,并且这种网络图十分庞大。 所以,Spark Graphx这种分布式计算框架十分适合这种场景。下面以设备间关联网络(节点数亿级别)为例,采用Graphx做一个设备团伙挖掘demo。团伙识别的算法采用的是Graphx自带的LabelPropagation算法。

下面的是Graphx示例代码(仅仅是demo):

其中输入文件格式:

A B weight

备注(A,B 代表设备id,String类型,weight:int,关联代表权重)

因为Graphx节点类型只支持Long,不支持String,所以,需要进行相应的转换,这里用到的广播变量进行idmap。

github链接: https://github.com/dylan-fan/spark_graphx_community_detection

[java] view plain copy
  1. package com.org.test  
  2.   
  3. import org.apache.spark.SparkConf  
  4. import org.apache.spark.SparkContext  
  5. import org.apache.spark.rdd.RDD  
  6. import org.apache.spark.graphx._  
  7. import scala.collection.mutable.Set  
  8.   
  9. object DeviceCom {  
  10.   def main(args: Array[String]) {  
  11.     if (args.length < 3) {  
  12.       println("usage: spark-submit com.org.test.DeviceCom <input> <output> <iternum>")  
  13.       System.exit(1)  
  14.     }  
  15.     val conf = new SparkConf()  
  16.     conf.setAppName("DeviceCom-" + System.getenv("USER"))  
  17.   
  18.     val sc = new SparkContext(conf)  
  19.   
  20.     val input = args(0)  
  21.   
  22.     val output = args(1)  
  23.   
  24.     val iternum = args(2).toInt  
  25.   
  26.     val vids = sc.textFile(input)  
  27.       .flatMap(line => line.split("\t").take(2))  
  28.       .distinct  
  29.       .zipWithUniqueId()  
  30.       .map(x => (x._1, x._2.toLong))  
  31.   
  32.     val vids_map = sc.broadcast(vids.collectAsMap())  
  33.   
  34.     val vids_rdd = vids.map {  
  35.       case (username, userid) =>  
  36.         (userid, username)  
  37.     }  
  38.   
  39.     val raw_edge = sc.textFile(input)  
  40.       .map(line => line.split("\t"))  
  41.     val col = raw_edge.collect()  
  42.   
  43.     val edges_rdd = sc.parallelize(col.map {  
  44.       case (x) =>  
  45.         (vids_map.value(x(0)), vids_map.value(x(1)))  
  46.     })  
  47.   
  48.     val g = Graph.fromEdgeTuples(edges_rdd, 1)  
  49.     val lp = lib.LabelPropagation.run(g, iternum).vertices  
  50.     val LpByUsername = vids_rdd.join(lp).map {  
  51.       case (id, (username, label)) =>  
  52.         (username, label)  
  53.     }  
  54.   
  55.     LpByUsername.map(x => x._1 + "\t" + x._2).saveAsTextFile(output)  
  56.   
  57.     sc.stop()  
  58.   }  
  59. }  


这里,只是采用Graphx做个demo(很简单啦),来测试Graphx在当前数据量级下的相关性能。实际设备团伙挖掘会更复杂,涉及到各种策略制定。

这篇关于###spark版### Spark Graphx 进行团伙的识别(community detection)的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/967677

相关文章

Nginx中配置使用非默认80端口进行服务的完整指南

《Nginx中配置使用非默认80端口进行服务的完整指南》在实际生产环境中,我们经常需要将Nginx配置在其他端口上运行,本文将详细介绍如何在Nginx中配置使用非默认端口进行服务,希望对大家有所帮助... 目录一、为什么需要使用非默认端口二、配置Nginx使用非默认端口的基本方法2.1 修改listen指令

MySQL按时间维度对亿级数据表进行平滑分表

《MySQL按时间维度对亿级数据表进行平滑分表》本文将以一个真实的4亿数据表分表案例为基础,详细介绍如何在不影响线上业务的情况下,完成按时间维度分表的完整过程,感兴趣的小伙伴可以了解一下... 目录引言一、为什么我们需要分表1.1 单表数据量过大的问题1.2 分表方案选型二、分表前的准备工作2.1 数据评估

MySQL进行分片合并的实现步骤

《MySQL进行分片合并的实现步骤》分片合并是指在分布式数据库系统中,将不同分片上的查询结果进行整合,以获得完整的查询结果,下面就来具体介绍一下,感兴趣的可以了解一下... 目录环境准备项目依赖数据源配置分片上下文分片查询和合并代码实现1. 查询单条记录2. 跨分片查询和合并测试结论分片合并(Shardin

SpringBoot结合Knife4j进行API分组授权管理配置详解

《SpringBoot结合Knife4j进行API分组授权管理配置详解》在现代的微服务架构中,API文档和授权管理是不可或缺的一部分,本文将介绍如何在SpringBoot应用中集成Knife4j,并进... 目录环境准备配置 Swagger配置 Swagger OpenAPI自定义 Swagger UI 底

基于Python Playwright进行前端性能测试的脚本实现

《基于PythonPlaywright进行前端性能测试的脚本实现》在当今Web应用开发中,性能优化是提升用户体验的关键因素之一,本文将介绍如何使用Playwright构建一个自动化性能测试工具,希望... 目录引言工具概述整体架构核心实现解析1. 浏览器初始化2. 性能数据收集3. 资源分析4. 关键性能指

Nginx进行平滑升级的实战指南(不中断服务版本更新)

《Nginx进行平滑升级的实战指南(不中断服务版本更新)》Nginx的平滑升级(也称为热升级)是一种在不停止服务的情况下更新Nginx版本或添加模块的方法,这种升级方式确保了服务的高可用性,避免了因升... 目录一.下载并编译新版Nginx1.下载解压2.编译二.替换可执行文件,并平滑升级1.替换可执行文件

Python进行JSON和Excel文件转换处理指南

《Python进行JSON和Excel文件转换处理指南》在数据交换与系统集成中,JSON与Excel是两种极为常见的数据格式,本文将介绍如何使用Python实现将JSON转换为格式化的Excel文件,... 目录将 jsON 导入为格式化 Excel将 Excel 导出为结构化 JSON处理嵌套 JSON:

一文解密Python进行监控进程的黑科技

《一文解密Python进行监控进程的黑科技》在计算机系统管理和应用性能优化中,监控进程的CPU、内存和IO使用率是非常重要的任务,下面我们就来讲讲如何Python写一个简单使用的监控进程的工具吧... 目录准备工作监控CPU使用率监控内存使用率监控IO使用率小工具代码整合在计算机系统管理和应用性能优化中,监

如何使用Lombok进行spring 注入

《如何使用Lombok进行spring注入》本文介绍如何用Lombok简化Spring注入,推荐优先使用setter注入,通过注解自动生成getter/setter及构造器,减少冗余代码,提升开发效... Lombok为了开发环境简化代码,好处不用多说。spring 注入方式为2种,构造器注入和setter

MySQL进行数据库审计的详细步骤和示例代码

《MySQL进行数据库审计的详细步骤和示例代码》数据库审计通过触发器、内置功能及第三方工具记录和监控数据库活动,确保安全、完整与合规,Java代码实现自动化日志记录,整合分析系统提升监控效率,本文给大... 目录一、数据库审计的基本概念二、使用触发器进行数据库审计1. 创建审计表2. 创建触发器三、Java