一、启动容器

# Docker 启动: docker run -d --restart=always --name neo4j -p 7474:7474 -p 7687:7687 -v /data/neo4j/data:/data -v /data/neo4j/import:/import -v /data/neo4j/logs:/logs -v /data/neo4j/conf:/conf -v /data/neo4j/plugins:/plugins neo4j:4.4.32

二、导入数据

LOAD CSV WITH HEADERS FROM 'file:///xueyuan.csv' AS line
MERGE (startNode:Table {name: line.from_table, type:line.from_type})
MERGE (endNode:Table {name: line.insert_table, type:line.insert_type})
MERGE (startNode)-[r:RELATES_TO {name: line.job_name, script_name: line.script_name, weight: line.weight, script_type:line.script_type}]->(endNode) 


-- 使用不同的标签或关系类型区分数据集
LOAD CSV WITH HEADERS FROM 'file:///xueyuan.csv' AS line
MERGE (insertNode:Table_New {name: line.insert_table})
MERGE (fromNode:Table_New {name: line.from_table})
MERGE (fromNode)-[r:RELATES_TO_NEW {type: line.use_name}]->(insertNode)



-- 导入job与job之间的依赖

LOAD CSV WITH HEADERS FROM 'file:///job_xueyuan.csv' AS line
MERGE (startNode:Job {name: line.start_job, job_type:line.start_job_type , cluster_name:line.start_cluster_name})
MERGE (endNode:Job {name: line.end_job, job_type:line.end_job_type , cluster_name:line.end_cluster_name})
MERGE (startNode)-[r:JOB_LIST {name: line.table_list}]->(endNode) 



LOAD CSV WITH HEADERS FROM 'file:///daobao8000.csv' AS line
MERGE (node1:Table {name: line.dev_sn_1})
MERGE (node2:Table {name: line.dev_sn_2})
WITH node1, node2, line
MERGE (node1)-[r:RELATES_TO]->(node2)
ON CREATE SET r.type = line.battery_id, r.weight = toFloat(line.w1)
MERGE (node2)-[r2:RELATES_TO]->(node1)
ON CREATE SET r2.type = line.battery_id, r2.weight = toFloat(line.w2)

三、查询语句

3.1 官方自带

查询所有层,注意增加层数,防止陷入循环遍历

  1. 查询下游

// 查询表的下游(一层)

// 查询表的下游(一层)
MATCH (startNode:Table {name: 'dwd.xxx'})-[r:RELATES_TO]->(endNode:Table)
RETURN startNode, r, endNode

//查询job 下游
MATCH (startNode:Job {name: 'xxxx'})-[r:JOB_LIST]->(endNode:Job)
RETURN startNode, r, endNode

// 增加过滤条件
MATCH (startNode:Table {name: 'xxxx'})-[r:RELATES_TO]->(endNode:Table)
WHERE NOT startNode.name CONTAINS 'temp' AND NOT endNode.name CONTAINS 'temp'
RETURN startNode, r, endNode


// 查询表的下游(所有层)r:RELATES_TO*..8 代表8层
MATCH (startNode:Table {name: 'dwd.xxxx'})-[r:RELATES_TO*..8]->(endNode:Table)
RETURN startNode, r, endNode
  1. 查询上游

// 查询表的上游(一层)
MATCH (startNode:Table)-[r:RELATES_TO]->(endNode:Table {name: 'dwd.xxxxx'})
RETURN startNode, r, endNode

MATCH (startNode:Table)-[r:RELATES_TO]->(endNode:Table {name: 'dwd.xxxxx'})
RETURN startNode, r, endNode

// 查询表的上游(所有层)r:RELATES_TO*..8 代表8层
MATCH path=(startNode:Table)-[r:RELATES_TO*..8]->(endNode:Table {name: 'dwd.xxxxxx'})
RETURN path

MATCH (startNode:Table )-[r:RELATES_TO*..8]->(endNode:Table{name: 'dwd.xxxxx'})
RETURN startNode, r, endNode
  1. 查询数据流

// 从sss.xxx表到dwd.xxxx表的数据流
MATCH path=(startNode:Table {name: 'sss.xxx'})-[:RELATES_TO*..5]->(endNode:Table {name: 'dwd.xxxx'})
RETURN path

MATCH path=(startNode:Table {name: 'sss.xxx'})-[:RELATES_TO*..5]->(endNode:Table {name: 'dwd.xxxx'})
RETURN path
  1. 查询所有的上游

MATCH (startNode:Table)-[r:RELATES_TO*..8]->(endNode:Table{name: 'ads.xxxxx'})
UNWIND r AS rel
RETURN DISTINCT rel.name AS relationship_name
  1. 查询多表的上游作业名

WITH [
    'dwd.t1',
    'dws.t2'
 
] AS tableNames
UNWIND tableNames AS tableName
MATCH (startNode:Table)-[r:RELATES_TO*..8]->(endNode:Table{name: tableName})
UNWIND r AS rel
RETURN DISTINCT tableName, rel.name AS relationship_name

3.2 使用APOC库(性能更好,推荐)

APOC提供了数据集成,数据导出,数据结构,高级图查询等诸多功能,本小节选取部分过程和函数进行演示。相比于过程,函数更容易理解,函数可以直接应用在Cypher查询中,对传入函数中的数据进行计算并返回计算后的结果,这点与Cypher内置的函数没有明显区别。过程的调用必须使用CALL命令,APOC中的过程可以类比与关系数据库中的存储过程。

  1. 查询作业、Node、Table下游

// 作业
MATCH (startNode:Job {name: 'xxxxx-yuancheng'})
CALL apoc.path.expandConfig(startNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: 'JOB_LIST>', // '>' 方向表示向下的路径
    labelFilter: '+Job',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, startNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
RETURN pathNodes AS endNodes, pathRels AS relationships, startNode;


// Table
MATCH (startNode:Table {name: 'dwd.xxxxxx'})
CALL apoc.path.expandConfig(startNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: 'RELATES_TO>', // '>' 方向表示向下的路径
    labelFilter: '+Table',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, startNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
RETURN pathNodes AS endNodes, pathRels AS relationships, startNode;


//通过表名查询下游所有的作业名
MATCH (startNode:Table {name: 'dwd.xxxxx'})
CALL apoc.path.expandConfig(startNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: 'RELATES_TO>', // '>' 方向表示向下的路径
    labelFilter: '+Table',               // 只包含 Table 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, startNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
UNWIND pathRels AS rel
RETURN distinct rel.name,rel.script_name,rel.owner;
  1. 查询上游


// 使用apoc插件后的查询语句 查询Job之间的依赖
MATCH (endNode:Job {name: 'xxxxxx'})
CALL apoc.path.expandConfig(endNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: '<JOB_LIST', // '<' 方向表示向上的路径
    labelFilter: '+Job',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, endNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
RETURN pathNodes AS startNodes, pathRels AS relationships, endNode;

//过滤条件
MATCH (endNode:Job {name: 'xxxxxxx'})
CALL apoc.path.expandConfig(endNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: '<JOB_LIST',   // '<' 方向表示向上的路径
    labelFilter: '+Job',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, endNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
WHERE ALL(node IN pathNodes WHERE node.job_type <> 'DRS' AND node.job_type <> 'FLINK')
RETURN pathNodes AS startNodes, pathRels AS relationships, endNode;




 // 使用apoc插件后的查询语句 查询 table
MATCH (endNode:Table {name: 'ads.xxxxx'})
CALL apoc.path.expandConfig(endNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: '<RELATES_TO', // '<' 方向表示向上的路径
    labelFilter: '+Table',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, endNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
//UNWIND pathNodes AS n
//RETURN distinct n.name;
RETURN pathNodes AS startNodes, pathRels AS relationships, endNode



查询上下游
MATCH (endNode:Table {name: 'data.xxxxx'})
CALL apoc.path.spanningTree(endNode, {
    relationshipFilter: '<RELATES_TO>', 
    labelFilter: '+Table', 
    maxLevel: -1,                       
    direction: 'BOTH'
})
YIELD path
WITH path, endNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
WHERE NONE(node IN pathNodes WHERE node.name CONTAINS 'cccc' OR node.name CONTAINS 'aaa' OR node.name CONTAINS 'bbbb' OR node.name CONTAINS 'sdj' OR node.name CONTAINS 'aghaj')
RETURN pathNodes AS startNodes, pathRels AS relationships, endNode;
  1. 查询数据流


MATCH (startNode:Job {name: 'big'}), (endNode:Job {name: 'dwd'})
CALL apoc.path.expandConfig(endNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: '<JOB_LIST', // '<' 方向表示向上的路径
    labelFilter: '+Job',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH path, endNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
UNWIND pathNodes AS startNodes
RETURN pathNodes AS startNodes, pathRels AS relationships, endNode;
  1. 查询表上游所有的作业名

// 通过Job级别的依赖进行查询,返回结果为当前节点与当前节点索引+1的节点
MATCH (endNode:Job {name: 'xxx'})
CALL apoc.path.expandConfig(endNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: '<JOB_LIST', // '<' 方向表示向上的路径
    labelFilter: '+Job',               // 只包含 Job 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
WITH endNode, nodes(path) AS pathNodes, relationships(path) AS pathRels
UNWIND range(0, size(pathNodes) - 1) AS idx  // 从0开始
WITH pathNodes[idx] AS currentNode, 
     pathNodes[idx +1] AS previousNode
RETURN DISTINCT currentNode.name AS currentNodeNamePart,
                previousNode.name AS previousNodeNamePart;
                
  1. 查询多元素的上游

-- 批量查询 作业的所有依赖与当前节点的上游节点
WITH [
    'j1'
] AS tableNames
UNWIND tableNames AS tableName
MATCH (startNode:Job {name: tableName})
CALL apoc.path.expandConfig(startNode, {
    uniqueness: 'NODE_GLOBAL',
    relationshipFilter: '<JOB_LIST', // '<' 方向表示向上的路径
    labelFilter: '+Job',
    maxLevel: -1
})
YIELD path
WITH startNode AS startTable, nodes(path) AS pathNodes
UNWIND pathNodes AS node
WITH distinct startTable, node AS upstreamTable
// 查询每个upstreamTable的上一层节点,并确保即使没有上游节点,也能返回一行数据
OPTIONAL MATCH (previousNode:Job)-[:JOB_LIST]->(upstreamNode:Job {name: upstreamTable.name})
WITH startTable, upstreamTable, collect(distinct previousNode) AS previousNodes
// 使用UNWIND展开previousNodes中的每个节点
UNWIND previousNodes AS previousNode
// 返回展开后的结果
RETURN  startTable.name as P0_job_name,
        startTable.start_time as P0_job_start_time,
        startTable.end_time as P0_job_end_time,
        upstreamTable.name AS current_job_name,
        upstreamTable.job_type AS current_job_type,
        upstreamTable.cluster_name AS current_cluster_name,
        upstreamTable.start_time as current_job_start_time,
        upstreamTable.end_time as current_job_end_time,
        previousNode.name AS previous_job_name,
        previousNode.job_type AS previous_job_type,
        previousNode.cluster_name AS previous_cluster_name,
        previousNode.start_time as previous_job_start_time,
        previousNode.end_time as previous_job_end_time;



WITH [
'j1',
'j2'
] AS tableNames
UNWIND tableNames AS tableName
MATCH (endNode:Job {name: tableName})
CALL apoc.path.expandConfig(endNode, {
    uniqueness: 'NODE_GLOBAL',        // 设置全局唯一性,防止重复访问节点
    relationshipFilter: '<JOB_LIST',  // '<' 方向表示向上的路径
    labelFilter: '+Job',              // 只包含 Job 标签的节点
    maxLevel: -1                      // 设置为-1表示无限制
})
YIELD path
WITH endNode, nodes(path) AS pathNodes
UNWIND pathNodes AS startNode
RETURN distinct endNode.name as hour_job_name,
        startNode.name AS job_name,
        startNode.owner as owner,
        startNode.start_time as job_start_time,
        startNode.end_time as job_end_timee,
        startNode.job_type as job_type




// 查询多张表的下游表名
WITH [
    't1',
    't2'
] AS tableNames
UNWIND tableNames AS tableName
MATCH (startNode:Table {name: tableName})
CALL apoc.path.expandConfig(startNode, {
    uniqueness: 'NODE_GLOBAL',         // 设置全局唯一性,防止重复访问节点
    relationshipFilter: 'RELATES_TO>', // '>' 方向表示向下的路径
    labelFilter: '+Table',             // 只包含 Table 标签的节点
    maxLevel: -1                       // 设置为-1表示无限制
})
YIELD path
UNWIND nodes(path) AS node
RETURN DISTINCT startNode.name AS startNodeName, node.name AS endNodeName;

更多推荐