Systems and methods for dynamic mapping for locality and balance
17 claims: 6 independent, 11 dependent
- 1コンピュータ実装方法であって、 コンピュータ・システムが、第1のパーティションのノードのヒストグラムを演算するステップと、 前記コンピュータ・システムが、第2のパーティションのノードのヒストグラムを演算するステップと、 前記コンピュータ・システムが、前記第1のパーティションのノードのヒストグラムに基づいて、前記第1のパーティションの一組のノードの候補パーティションとして前記第2のパーティションを選択するステップと、 前記コンピュータ・システムが、前記第2のパーティションのノードのヒストグラムに基づいて、前記第2のパーティションの一組のノードの候補パーティションとして前記第1のパーティションを選択するステップと、 前記コンピュータ・システムが、負荷平衡に基づいて、前記第1のパーティションの一組のノードの少なくとも一部を前記第2のパーティションに再マッピングし、前記第2のパーティションの一組のノードの少なくとも一部を前記第1のパーティションに再マッピングするステップと、を備える、コンピュータ実装方法。
- 2前記コンピュータ・システムが、第3のパーティションのノードのヒストグラムを演算するステップと、 前記コンピュータ・システムが、前記第1のパーティションのノードのヒストグラムに基づいて、前記第1のパーティションの別の一組のノードの候補パーティションとして前記第3のパーティションを選択するステップと、 前記コンピュータ・システムが、前記第3のパーティションのノードのヒストグラムに基づいて、前記第3のパーティションの一組のノードの候補パーティションとして前記第1のパーティションを選択するステップと、 前記コンピュータ・システムが、負荷平衡に基づいて、前記第1のパーティションの他の一組のノードの少なくとも一部を前記第3のパーティションに再マッピングするとともに、前記第3のパーティションの一組のノードの少なくとも一部を前記第1のパーティションに再マッピングするステップと、をさらに備える、請求項1に記載のコンピュータ実装方法。
- 3前記コンピュータ・システムが、エッジ局所性のゲインに基づいて、前記第1のパーティションの一組のノードをソートするステップと、 前記コンピュータ・システムが、エッジ局所性のゲインに基づいて、前記第2のパーティションの一組のノードをソートするステップと、をさらに備え、 前記第2のパーティションが、エッジ局所性のゲインに関する確率に基づいて、前記第1のパーティションのノードの候補パーティションとして選択される、請求項1または2に記載のコンピュータ実装方法。
- 4前記第1のパーティションのノードのヒストグラムが、複数のパーティションのそれぞれにおける接続ノードの数を示す、請求項1~3のいずれか一項に記載のコンピュータ実装方法。
- 5前記第2のパーティションに再マッピングされた前記第1のパーティションのノードの数と前記第1のパーティションに再マッピングされた前記第2のパーティションのノードの数との差が、しきい値の範囲内であること、および、 前記第2のパーティションに再マッピングされた前記第1のパーティションのノードの重みと前記第1のパーティションに再マッピングされた前記第2のパーティションのノードの重みとの差が、しきい値の範囲内であることのうちの少なくとも一方を含む、請求項1~4のいずれか一項に記載のコンピュータ実装方法。
- 6前記コンピュータ・システムが、前記再マッピングの前に、前記第1のパーティションの第1の総ノード重みを演算するステップをさらに備え、 好ましくは、前記コンピュータ・システムが、前記再マッピングの後に、前記第1のパーティションの第2の総ノード重みを演算するステップをさらに備える、請求項1~5のいずれか一項に記載のコンピュータ実装方法。
- 7前記コンピュータ・システムが、非分散システムであり、 該コンピュータ実装方法が、 前記コンピュータ・システムが、ノードグラフをメモリにロードするステップをさらに備え、 前記ノードグラフが、前記第1のパーティションのノードおよび前記第2のパーティションのノードを含む、請求項1~6のいずれか一項に記載のコンピュータ実装方法。
- 8前記コンピュータ・システムが、分散システムであり、 該コンピュータ実装方法が、 前記コンピュータ・システムが、ノードグラフの異なる部分を前記分散システム全体にロードするステップをさらに備え、 前記ノードグラフが、前記第1のパーティションのノードおよび前記第2のパーティションのノードを含む、請求項1~6のいずれか一項に記載のコンピュータ実装方法。
- 9前記第1のパーティションの複数のノードのそれぞれと関連付けられた接続ノードの現行パーティションIDを受信するステップをさらに備え、 好ましくは、前記第1のパーティションのノードのヒストグラムが、前記現行パーティションIDに基づいて演算され、 好ましくは、前記第1のパーティションの複数のノードのそれぞれの現行パーティションIDを提供するステップをさらに備える、請求項1~8のいずれか一項に記載のコンピュータ実装方法。
- 10候補パーティションが、局所性ゲインしきい値に基づいて選択される、請求項1~9のいずれか一項に記載のコンピュータ実装方法。
- 11前記第2のパーティションが、エッジ局所性のゲインに関する確率に基づいて、前記第1のパーティションのノードの候補パーティションとして選択される、請求項1~10のいずれか一項に記載のコンピュータ実装方法。
- 12前記コンピュータ・システムが、再マッピングされるノードを示す複数のパーティションのすべてのパーティション対の記録を生成するステップをさらに備える、請求項1~11のいずれか一項に記載のコンピュータ実装方法。
- 13前記ノードグラフが、ソーシャルネットワーキング・システムによりサポートされている、請求項1~12のいずれか一項に記載のコンピュータ実装方法。
- 14少なくとも1つのプロセッサと、 前記少なくとも1つのプロセッサに指示して、請求項1~13のいずれか一項に記載の方法を実行するように構成された命令を格納したメモリと、 を備えた、システム。
- 15実行された場合に、請求項1~13のいずれか一項に記載のコンピュータ実装方法をコンピュータ・システムに実行させるコンピュータ実行可能命令を格納した、コンピュータ記憶媒体。
- 16システムであって、 少なくとも1つのプロセッサと、 メモリと、を備え、 前記メモリは、 前記少なくとも1つのプロセッサに指示して、 第1のパーティションのノードのヒストグラムを演算すること、 第2のパーティションのノードのヒストグラムを演算すること、 前記第1のパーティションのノードのヒストグラムに基づいて、前記第1のパーティションの一組のノードの候補パーティションとして前記第2のパーティションを選択すること、 前記第2のパーティションのノードのヒストグラムに基づいて、前記第2のパーティションの一組のノードの候補パーティションとして前記第1のパーティションを選択すること、 負荷平衡に基づいて、前記第1のパーティションの前記一組のノードの少なくとも一部を前記第2のパーティションに再マッピングするとともに、前記第2のパーティションの前記一組のノードの少なくとも一部を前記第1のパーティションに再マッピングすることと、を実行するように構成された命令を格納する、システム。
- 17コンピュータ記憶媒体であって、 実行された場合に、 第1のパーティションのノードのヒストグラムを演算するステップと、 第2のパーティションのノードのヒストグラムを演算するステップと、 前記第1のパーティションのノードのヒストグラムに基づいて、前記第1のパーティションの一組のノードの候補パーティションとして前記第2のパーティションを選択するステップと、 前記第2のパーティションのノードのヒストグラムに基づいて、前記第2のパーティションの一組のノードの候補パーティションとして前記第1のパーティションを選択するステップと、 負荷平衡に基づいて、前記第1のパーティションの一組のノードの少なくとも一部を前記第2のパーティションに再マッピングするとともに、前記第2のパーティションの一組のノードの少なくとも一部を前記第1のパーティションに再マッピングするステップと、 を含むコンピュータ実装方法をコンピュータ・システムに実行させるコンピュータ実行可能命令を格納した、コンピュータ記憶媒体。
Independent claims17
94 paragraphs, as filed
0001The present invention relates to the field of node graphs, in particular computer mounting methods and systems as well as storage media. More specifically, the present invention provides techniques for mapping nodes to partitions.
0002Social networking websites provide a dynamic environment in which members can connect and communicate with other members. These websites may generally provide an online mechanism that allows members to interact in existing social networks and create new social networks. Members may include any individual or entity, such as an organization or company. Among the many attributes, social networking websites enable members to effectively and efficiently communicate relevant information to their social networks.
0003Members of social networks emphasize or share information, news articles, relationship activities, music, video, and any other content of interest to areas of the website that are available for content exclusively or otherwise. be able to. Other members of the social network can access the shared content by viewing the member profile or performing a dedicated search. In accessing and evaluating the content, other members may respond with one or more response actions, such as providing feedback or comment on the content. This ability for members to interact with each other facilitates communication between members and helps achieve the goals of social networking websites.
0004The social network may be modeled as a social graph. A node graph, such as a social graph, may contain a huge number of nodes and edges connecting the nodes. In the case of social networking systems, users can access and share vast amounts of information reflected in node graphs. For example, the number of nodes can be in the hundreds of millions or billions. Maintaining and providing such a large amount of data involves many challenges.
<p num="0005"> To dynamically map nodes for locality and equilibrium, computer implementations, systems, and computer-readable media may, in one embodiment, compute a histogram of the nodes in the first partition. In addition, the histogram of the node of the second partition may be calculated. The second partition may be selected as a candidate partition for a set of nodes in the first partition based on the histogram of the nodes in the first partition. The first partition may be selected as a candidate partition for a set of nodes in the second partition based on the histogram of the nodes in the second partition. Based on load balancing, at least part of the set of nodes in the first partition is mapped to the second partition, and at least part of the set of nodes in the second partition is mapped to the first partition. It may be done.</p><p num="0006"> In one embodiment, the histogram of the nodes of the third partition may be calculated. The third partition may be selected as a candidate partition for another set of nodes in the first partition based on the histogram of the nodes in the first partition. The first partition may be selected as a candidate partition for a set of nodes in the third partition based on the histogram of the nodes in the third partition. Based on load balancing, at least a portion of the other set of nodes in the first partition is mapped to the third partition, and at least a portion of the set of nodes in the third partition is the first. It may be mapped to a partition.</p><p num="0007"> In one embodiment, the set of nodes in the first partition may be sorted based on the gain of edge locality. A set of nodes in the second partition may be sorted based on the gain of edge locality.</p><p num="0008"> In one embodiment, the second partition may be selected as a candidate partition for the node of the first partition based on the probability of gain of edge locality.</p><p num="0009"> In one embodiment, the histogram of the nodes in the first partition may show the number of connected nodes in each of the plurality of partitions. In one embodiment, the difference between the number of nodes in the first partition remapped to the second partition and the number of nodes in the second partition remapped to the first partition is within the threshold range. It may be inside.</p><p num="0010"> In one embodiment, the difference between the node weights of the first partition remapped to the second partition and the node weights of the second partition remapped to the first partition is within the threshold range. It may be inside.</p><p num="0011"> In one embodiment, the first total node weight of the first partition may be calculated prior to remapping. In one embodiment, the second total node weight of the first partition may be calculated after the remapping.</p><p num="0012"> In one embodiment, the computer system may be a non-distributed system. Also, the node graph may be loaded into memory. The node graph may include nodes in the first partition and nodes in the second partition.</p><p num="0013"> In one embodiment, the computer system may be a distributed system. Also, different parts of the node graph may be loaded throughout the distributed system. The node graph may include nodes in the first partition and nodes in the second partition.</p><p num="0014"> In one embodiment, the current partition ID of the connection node associated with each node of the first partition may be received. In one embodiment, the histogram of the nodes in the first partition is calculated based on the current partition ID.</p><p num="0015"> In one embodiment, the current partition ID of each node of the first partition may be provided. In one embodiment, candidate partitions may be selected based on the local gain threshold.</p><p num="0016"> In one embodiment, the second partition may be selected as a candidate partition for the node of the first partition based on the probability of gain of edge locality.</p><p num="0017"> In one embodiment, a record of all partition pairs of multiple partitions indicating the nodes to be remapped may be generated. In one embodiment, the node graph may be supported by a social networking system.</p><p num="0018"> The accompanying drawings and the detailed description below will reveal many other features and embodiments of the invention. The embodiments according to the present invention are specifically disclosed in the appended claims relating to methods, systems, and media, and any feature described in one claim category (eg, method) may be any feature in another claim category (eg, method). It can be billed in the system) as well.</p><p num="0019"> In one embodiment of the present invention, the computer mounting method is The steps of computing the histogram of the nodes in the first partition by the computer system, The steps of computing the histogram of the nodes in the second partition by the computer system, A step in which the computer system selects the second partition as a candidate partition for a set of nodes in the first partition based on the histogram of the nodes in the first partition. A step in which the computer system selects the first partition as a candidate partition for a set of nodes in the second partition based on the histogram of the nodes in the second partition. The computer system remaps at least a portion of the set of nodes in the first partition to the second partition, and at least a portion of the set of nodes in the second partition, based on load balancing. Steps to remap to partition 1 and including.</p><p num="0020"> This computer implementation method The steps of computing the histogram of the nodes in the third partition by the computer system, A step in which the computer system selects a third partition as a candidate partition for another set of nodes in the first partition, based on the histogram of the nodes in the first partition. A step in which the computer system selects the first partition as a candidate partition for a set of nodes in the third partition based on the histogram of the nodes in the third partition. The computer system remaps at least some of the other set of nodes in the first partition to the third partition based on load balancing, and at least one of the set of nodes in the third partition. Steps to remap the part to the first partition, It is preferable to further include.</p><p num="0021"> Also, this computer implementation method The steps by the computer system to sort a set of nodes in the first partition based on the gain of edge locality, A step in which the computer system sorts a set of nodes in the second partition based on the gain of edge locality, May include.</p><p num="0022"> In the above, the second partition is preferably selected as a candidate partition for the node of the first partition based on the probability of gain of edge locality.</p><p num="0023"> In another embodiment, the histogram of the nodes in the first partition shows the number of connected nodes in each of the plurality of partitions, and / or The difference between the number of nodes in the first partition remapped to the second partition and the number of nodes in the second partition remapped to the first partition is within the thresholds, and / Or The difference between the node weights of the first partition remapped to the second partition and the node weights of the second partition remapped to the first partition is within the threshold.</p><p num="0024"> In another embodiment, this method Before remapping, the step of calculating the first total node weight of the first partition, and After the remapping, it involves both or one of the steps of calculating the second total node weight of the first partition.</p><p num="0025"> Also, in this computer implementation method, the computer system is a non-distributed system. The computer system may further include a step of loading a node graph into memory, wherein the node graph includes a node in the first partition and a node in the second partition.</p><p num="0026"> Also, in this computer implementation method, the computer system is a distributed system, The step of loading different parts of the node graph into the entire distributed system by the computer system, the node graph may further include a step containing a node in the first partition and a node in the second partition. , And / or It further includes the step of receiving the current partition ID of the connection node associated with each of the nodes in the first partition, and / or The histogram of the nodes in the first partition is calculated based on the current partition ID and / or It further includes a step of providing the current partition ID for each node in the first partition, and / or Candidate partitions are selected based on the local gain threshold and / or The second partition is selected as a candidate partition for the node of the first partition based on the probability of edge locality gain and / or It further includes the step of generating a record of all partition pairs of multiple partitions indicating the nodes to be remapped by the computer system.</p><p num="0027"> This computer implementation method may also include node graphs supported by social networking systems. The system according to the present invention With at least one processor A memory containing instructions configured to instruct at least one processor to perform a method according to any or all of the above embodiments. It is preferable to provide.</p><p num="0028"> A computer storage medium containing computer-executable instructions that, when executed, causes a computer system to perform the computer implementation method according to any or all of the above embodiments.</p><p num="0029"> The system according to the present invention With at least one processor Instruct at least one processor above Computing the histogram of the nodes in the first partition, Computing the histogram of the nodes in the second partition, To select the second partition as a candidate partition for a set of nodes in the first partition, based on the histogram of the nodes in the first partition. To select the first partition as a candidate partition for a set of nodes in the second partition based on the histogram of the nodes in the second partition. Based on load balancing, remap at least some of the nodes in the first partition to the second partition, and remap at least some of the nodes in the second partition to the first partition. Mapping and And the memory that stores the instructions that are configured to execute It is preferable to provide.</p><p num="0030"> When executed, Steps to calculate the histogram of the nodes in the first partition, Steps to calculate the histogram of the nodes in the second partition, The step of selecting the second partition as a candidate partition for a set of nodes in the first partition based on the histogram of the nodes in the first partition, Steps to select the first partition as a candidate partition for a set of nodes in the second partition based on the histogram of the nodes in the second partition, Based on load balancing, remap at least some of the nodes in the first partition to the second partition, and remap at least some of the nodes in the second partition to the first partition. The steps to map and A computer storage medium containing computer-executable instructions that causes a computer system to execute computer implementation methods, preferably including.</p>
0031<figref num="1">It is a figure which showed the exemplary optimization module 100 which concerns on one Embodiment.</figref><figref num="2">It is a figure which showed the exemplary locality control module 102 which concerns on one Embodiment.</figref><figref num="3">It is a figure which showed the exemplary distributed system which concerns on one Embodiment.</figref><figref num="4">It is a figure which showed the example equilibrium remapping module 103 which concerns on one Embodiment.</figref><figref num="5">It is a figure which showed the exemplary optimization process which remaps a node which concerns on one Embodiment.</figref><figref num="6">FIG. 6 is an exemplary network diagram of a system that optimizes node mapping within a social networking system according to an embodiment.</figref><figref num="7">FIG. 5 illustrates an exemplary computer system that can be used in one or more implementations of the embodiments described herein according to an embodiment.</figref>
0032The drawings show various embodiments of the invention for illustration purposes only, and similar reference numerals are used to identify similar elements. From the following description, one of ordinary skill in the art will readily recognize that other embodiments of the structures and methods shown in the drawings can be employed without departing from the principles of the invention described herein.
0033The node graph may be distributed across multiple computers in a distributed system. How the nodes are mapped to partitions can affect performance. For example, in a social network, if you map a friend (for example, a node connected by an edge based on "friendship") to the same partition, a fanout query against another partition when a query about the friend is initiated Performance may be improved by minimizing the amount. The partition may contain, for example, a server or data center. In addition, friends tend to query for duplicate information, which improves network performance because it is likely that the duplicate information has already been read for the query and stored in high-speed memory such as a cache. there's a possibility that. Therefore, it may be beneficial to increase edge locality.
0034Furthermore, the partition load may also affect network performance. The type of "load" can vary depending on the application. For example, in some applications the load may be based on the number of nodes (eg users) in the partition. In other applications, the load may be based on the amount of data stored in the partition. In yet other applications, the load may be based on activity level. The activity level may be related to the user's activity on the partition. The activity may be related to the number of times the user performs one or more activities, the amount of data the user uploads, downloads, or both. The type of load can vary from application to application, but is generally beneficial to the ability to maintain a balanced load between partitions.
0035In some cases, the nodes may be mapped to partitions in a roughly random manner that is likely to be low edge locality. For example, users of a social network may be mapped to that partition based on their time of participation in the network and the partitions available at that time. In such cases, users of the social network may be able to map to partitions without considering where each friend is mapped.
0036In some cases, higher edge while balancing the load between partitions The initial mapping of the node graph may be generated when trying to provide locality). An exemplary initial mapping may be based, for example, on a geographic location that seeks to generate high edge locality while balancing loads between partitions. However, node graphs such as social graphs change dynamically and frequently, which can adversely affect edge locality or equilibrium. For example, users of social networks may make or lose friends, move to different geographic locations, and join new organizations, companies, schools, and so on. In addition, new users may join the network and form new friendships with existing users, which can also adversely affect edge locality or equilibrium. In addition, users may change their behavior or habits, which can change the amount of load they place on the mapped partitions. Some initial mapping methods may be useful at the time of creation, but they are static and may not support changing node graphs.
0037Embodiments of the systems and methods described herein relate to optimizing node mapping in node graphs to improve edge locality (or connection locality) while maintaining load equilibrium in partitions. The optimization process may include a series of iterations that constantly improve edge locality in the partition while still maintaining load balancing. For example, in certain embodiments, the optimization process may attempt to determine which nodes benefit edge locality when remapped to different partitions. The candidate partition may be selected for each node that enjoys the benefits of remapping to the candidate partition. Then, in the optimization process, the nodes may be selectively moved with respect to the candidate partition in such a manner as to maintain the load balance. The optimization process may then perform a series of subsequent iterations until the optimization process stabilizes.
0038In addition, the optimization process may be repeated in the future to optimize the mapping to accommodate any changes in the node graph that occur over time. In one embodiment, the optimization process may be re-executed after a predetermined period of time, such as after one week, one month, six months, or any other period.
0039If the node graph changes over time, these changes can be detrimental to edge locality or equilibrium, which can lead to discrepancies (or disturbing effects) on stable mapping. Since the other stable parts of the mapping are already stable, the optimization process primarily focuses on these discrepancies. By doing so, the optimization process may be able to constantly optimize the current mapping without having to generate an entirely new mapping from scratch each time. For example, new users tend to have no friendship connections, so they are initially mapped to random partitions with minimal or no edge locality. You may be. Future optimization processes may focus on specific reasons for change to stable mapping.
0040FIG. 1 shows an exemplary optimization module 100 according to an embodiment. The optimization module 100 optimizes the node mapping of the node graph to the partition so as to improve edge locality (or connection locality) while maintaining load equilibrium across the partition. In one embodiment, the nodes and edges described herein may relate to users of social networking systems and their respective friendships. Also, the underlying ideas and principles may be applicable to other types of nodes and edges. In one embodiment, a node, whether concrete or abstract, can be represented as a human, non-human, organization, content (eg, image, video, audio, etc.), event, web page, communication, object, concept. , Or any other thing, idea, construct, etc. The node may include users of networking systems such as social networking systems. The user is not necessarily limited to humans and may include other non-human entities. Edge locality may be in terms of the number or proportion of edges within a partition, as opposed to the edges between two partitions. For a given partition, edge locality may be with respect to the number of edges contained in the given partition, as opposed to the number of edges leading to different partitions. For a given node, edge locality refers to the number of edges in the partition of the given node connected to the given node, relative to the number of edges connecting the given node to different partitions. You may. The optimization module 100 includes an initialization module 101, a locality control module 102, and a balanced remapping module. module) 103 may be provided. The components shown in this drawing and all the drawings herein are merely exemplary, and in other embodiments, they may include additional components, fewer components, or different components. Some components are not shown so that the relevant details are not obscured.
0041The initialization module 101 may get the current mapping (eg, initial mapping) of the nodes in the node graph to the partition. The initialization module 101 may load the node graph into memory. The initialization module 101 may also load node weights or edge weights into memory. The node weight may be related to the amount of load the user has placed on the partition. For example, a user's node weight in a social networking system may be related to the number of user logins, photo uploads, and so on. Edge weights may relate to the "cost" of having two nodes in different partitions, such as the amount of load applied to the edge between two users. For example, two users who share a large amount of data with each other may be determined to have high edge weights. The edge factor may be related to the proximity or closeness of the two nodes. For example, two users may have a higher coefficient than a general acquaintance by being a family member. In some cases, the coefficients may act as edge weights.
0042The initialization module 101 may also load the initial mapping of the nodes into the partition. For example, the initial mapping may include an initial mapping of a node to a partition based on the geographic relevance of the node. Then, as described in the present specification, the locality may be improved while maintaining the load equilibrium by optimizing the initial mapping.
0043In one embodiment, a non-distributed system optimizes the mapping. For example, a node graph is loaded into memory (for example, random access memory (RAM)) on a single or local machine (or other non-distributed system) that can compute the optimized mapping of nodes to partitions. It may be like this.
0044In another embodiment, a distributed system with multiple computers may be adapted to optimize the mapping. For example, node graphs may be loaded across multiple computers (or machines) so that each computer loads a unique set of node graphs. The node graph may be included in a table by a distributed file system such as Hadoop Distributed File System (HDFS). Also, node weights, edge weights, and initial mappings may be contained in separate tables by the distributed file system.
0045The locality control module 102 performs arithmetic and analysis related to the improvement of edge locality in the partition, as described in more detail herein. Equilibrium remapping module 103 performs operations and analysis on node remapping to different partitions in a manner that maintains load equilibrium across the partitions, as described in more detail herein.
0046FIG. 2 shows an exemplary locality control module 102 according to an embodiment. The locality control module 102 includes a histogram generation module 201, a gain computation module 202, a candidate partition selection module 203, and a matrix generation module 204. It may be provided.
0047Histogram generation module 201 may compute the histogram of nodes in one or more partitions. The number of nodes (or nodes that have a connection to a given node) connected to a given node by an edge may be referred to herein as "connecting nodes". In one embodiment that does not consider edge weights, the histogram of a given node may be able to identify the number of connected nodes within each partition across multiple partitions. In one embodiment that considers edge weights, the histogram of a given node may be able to identify the total weight of the edges connected to each partition across multiple partitions.
0048In one embodiment where the non-distributed system optimizes the mapping, the histogram of each node may be calculated by the non-distributed system (eg, a single computer). For example, a single computer may contain data that identifies each node's connecting node, as well as the current partition of each node. In this way, the computer may access the data and generate a histogram of each node.
0049In one embodiment in which a distributed system optimizes mapping, the use of multiple computers may result in the use of multiple computers to generate a histogram of the nodes of a node graph by exchanging information (or messages) between the computers. .. For example, each node may propagate its current partition to each of its connection nodes (eg, via the current partition ID). Each node may then generate its own histogram after each connecting node identifies the mapped partition. In one embodiment that considers node weights, node weights may also be exchanged between nodes. As described in more detail herein, the total weights of the nodes or edges within each partition may be identified by aggregating the current partition and node weights into a table as a whole.
0050The gain calculation module 202 may be made to calculate the "gain" value associated with the remapping of the node from the current partition to a different partition. The "gain" of a given node is now calculated by the difference between the number of connected nodes mapped to different partitions and the number of connected nodes mapped to the current partition of a given node. Good. For example, if a given node is assigned to partition # 1, 5 connected nodes are mapped to partition # 1, and 7 connected nodes are mapped to partition # 2, then the given node. May have a gain of "2" (ie, 7-5) when remapped to partition # 2. The gain associated with remapping to a different partition may be positive, increasing the edge locality of a given node, reflecting an increase in the number of connected nodes in the same partition as a given node. Let me. The gain associated with remapping to a different partition can be negative and reduces the edge locality of a given node, reflecting a decrease in the number of connected nodes in the same partition as a given node. Let me. In some cases, the gain may be 0, reflecting the same number of connected nodes in the same partition as a given node.
0051Considering Edge Weights In one embodiment, the gain of a given node associated with remapping to different partitions is the total edge weights of the connecting nodes mapped to different partitions and the given node. It may be calculated by the difference from the total edge weight of the connection node mapped to the current partition of.
0052Candidate partition selection module 203 may choose a candidate partition to which the node is remapped. The selection of candidate partitions for a given node may be made based on the improvement in edge locality of the given node. In some cases, candidate partitions that maintain the same edge locality (that is, have a gain of 0) may be selected for a given node.
0053In one embodiment, candidate partitions for a given node may be selected based on the probability of gaining the edge locality of the given node. For example, candidate partitions for a given node may be selected based on a probabilistic algorithm that has a bias towards partitions with greater edge locality. Thus, a partition with a higher gain for a given node is more likely to be selected as a candidate partition for the given node, but it is not always certain. For example, node C is mapped to partition # 3, with partition # 3 having 5 connection nodes, partition # 6 having 10 connection nodes, partition # 12 having 8 connection nodes, and other partitions. , It does not have to have a connection node. In one embodiment, the probability algorithm may be defined based on the number of connected nodes in the partition as well as the total number of connected nodes in the partition for which gain is obtained. For example, the probability of partition # 3 being selected as a candidate partition is 5/23, where "5" is the number of connected nodes in partition # 3 and "23" is the total number of connected nodes (partition # 3). From 5 to 10, partition # 6 to 10 and partition # 12 to 8). The probability of partition # 6 being selected as a candidate partition is 10/23. The probability of partition # 12 selected as a candidate partition is 8/23. In one embodiment, a partition having a gain of 0 can also be considered as a candidate partition. For example, in this way, the probability algorithm may be defined based on the number of connected nodes in the partition, as well as the total number of connected nodes in the partition with a positive or zero gain.
0054In another embodiment, the partition with the maximum edge locality may be selected as the candidate partition. In one embodiment, a local gain threshold may be implemented that must meet or exceed the threshold to qualify the partition as a candidate partition. For example, if the local gain threshold is 3, then a partition with a gain of 1 or 2 for a given node remapping is not eligible for selection as a candidate partition. However, partitions with a gain of 3 or more for remapping a given node are eligible for selection as candidate partitions. For example, implementation of a local gain threshold may suppress or eliminate remapping between partitions with the lowest gain.
0055In one embodiment, the local gain threshold may change with each optimization iteration. For example, the local gain threshold may start at a high value and decrease over iteration. In addition, other change patterns may be implemented with respect to the threshold value. In some cases, high local gain thresholds may be set to stop or slow down node remapping over a long or short period of time. In one embodiment in which the distributed system optimizes the mapping, the local gain threshold may start at a high value and decrease over iteration. Thus, a more useful remapping of the node may occur first and then gradually decline as the distribution improves and stabilizes to provide a finer remapping. ..
0056The matrix generation module 204 generates a matrix of node lists based on the nodes of every partition pair (X, Y). In one embodiment where the non-distributed system optimizes the mapping, the list of partition pairs (X, Y) lists the nodes in partition X with partition Y as the candidate partition. For example, in the case of partition pair (1,2), all the nodes in partition # 1 with partition # 2 as a candidate partition may be listed. In the case of partition pair (1,3), all the nodes in partition # 1 with partition # 3 as a candidate partition may be listed. In the case of partition pair (3,1), all the nodes in partition # 3 with partition # 1 as a candidate partition may be listed. And a list of all partition pairs (X, Y) may be generated based on this data. Also, in one embodiment, the list may indicate the gain associated with the node that is remapped to the candidate partition. Considering Node Weights In one embodiment, the list of partition pairs (X, Y) may include the total weights of the nodes in partition X with partition Y as the candidate partition.
0057In one embodiment in which the distributed system optimizes the mapping, the list in the matrix may include a value representing the number of nodes in partition X with partition Y as the candidate partition. When considering node weights, the list may include the total node weights of the nodes of partition X with partition Y as the candidate partition. In one embodiment, the machine associated with each node may propagate its current partition and node weights to machines associated with other nodes in other partitions. In this way, the total node weight of each partition may be reflected by aggregating the node weight as a whole. Also, the total node weight of the partition may be reflected in the table.
0058The list generated by the matrix generation module 204 may then be provided to the balanced remapping module 103 to remap the nodes to different partitions. As described in more detail herein, remapping may be performed to increase edge locality while maintaining partition load balance.
0059The structure of the non-distributed system may differ from that of the distributed system. For example, in one embodiment in which a non-distributed system optimizes mapping, the histogram generation module 201, the gain calculation module 202, the candidate partition selection module 203, and the matrix generation module 204 are implemented in a single computer system. Often, the data associated with the optimization process may be held in a single or shared memory (eg, RAM) of the computer system. In a distributed system, the gain calculation module 202, the candidate partition selection module 203, and the matrix generation module 204 may be implemented in one or more computer systems of the distributed system, respectively.
0060FIG. 3 shows an exemplary distributed system according to an embodiment. The distributed system 300 may include computer systems 301-303, aggregation modules 304-312, and master computer system 313.
0061The node graphs may be loaded across computer systems 301-303 such that each of computer systems 301-303 contains a unique (or non-unique) node graph subset. As an example, three computer systems are used, but of course, in other embodiments, any number of other computer systems may be implemented.
0062Also, computer systems 301-303 may each include a subset of aggregation modules 304-312. Aggregation modules 304 to 312 may each generate a node list and a node list of one or more partition pairs (X, Y) that can be used to generate a matrix of partition pairs (X, Y). In one embodiment, each of the aggregation modules 304 to 312 may generate a list of one of a pair of partitions (X, Y). Further, the aggregation modules 304 to 312 may include, for example, a histogram generation module 201, a gain calculation module 202, and a candidate partition selection module 203, respectively. As an example, nine aggregation modules 304 to 312 are used, but of course, in other embodiments, the number of aggregation modules implemented may be other than this. Also, of course, the number of aggregation modules per computer system may be different.
0063Each computer system 301-303 may receive a node list from a subset of its respective aggregation modules 304-312. Computer systems 301-303 may then communicate each list to each other and to master computer system 313, respectively. For example, computer system 301 may propagate the list generated by aggregation modules 304-306 to computer systems 302 and 303 as well as master computer system 313. Computer systems 301-303 and master computer systems 313 may then generate the entire matrix of all partition pairs (X, Y). As described above, the matrix generation module 204 may be implemented by, for example, aggregation modules 304 to 312, computer systems 301 to 303, and master computer system 313. Aggregation modules 304-312 each generate a list of specific partition pairs (X, Y), while computer systems 301-303 and master computer systems 313 generate an entire matrix of partition pairs (X, Y). It may be designed to do. In one embodiment, the master computer system 313 produces a node list and a matrix of partition pairs (X, Y), although it does not need to include any subset of the node graph.
0064FIG. 4 shows an exemplary equilibrium remapping module 103 according to an embodiment. The balanced remapping module 103 may include a node selection module 401, a node remapping module 402, and a termination module 403.
0065The node selection module 401 can remap which node of the first partition to the second partition for any two given partitions in order to increase edge locality while maintaining load balance across the partitions. You may decide whether or not. Node remapping module 402 remaps nodes between any given two partitions determined by node selection module 401.
0066In one embodiment that does not consider node weights, the total number of nodes remapped from the first partition to the second partition is the same as the total number of nodes remapped from the second partition to the first partition. You may. In another embodiment, the total number of nodes remapped from the first partition to the second partition is not the same as the total number of nodes remapped from the second partition to the first partition, but the entire partition. It may be within the permissible range for maintaining load equilibrium. Considering Node Weights In one embodiment, the total weight of the nodes remapped from the first partition to the second partition is the same as the total weight of the nodes remapped from the second partition to the first partition. It may be. In another embodiment, the total weight of the nodes remapped from the first partition to the second partition is not the same as the total weight of the nodes remapped from the second partition to the first partition, It may be within the permissible range considered to be load equilibrium.
0067The node selection module 401 may determine the remapping based on the partition pair (X, Y) generation list. In one embodiment, to maintain equilibrium, the number of nodes selected to move from partition X to partition Y may be the same as the number of nodes selected to move from partition Y to partition X. .. If the list of two partition pairs (eg, partition pair (1,2) and partition pair (2,1)) has different numbers of nodes, in one embodiment, remapping from each partition is selected for the node. The number may be the same as the number of nodes in the list of partition pairs with the smaller number of nodes. For example, the list of partition pairs (1,2) contains 5 nodes of partition # 1 with partition # 2 as candidate partitions, and 10 nodes of partition # 2 with partition # 1 as candidate partitions (partition pairs (1,2). If the list of 2,1) is included, 5 nodes may be selected for remapping from each partition. For example, all five nodes in the partition pair (1,2) list are remapped to partition # 2, and the five nodes in the partition pair (2,1) list with the highest gain are remapped to partition # 1. It may be mapped.
0068In one embodiment where the non-distributed system optimizes the mapping, the node selection module 401 may include a sort module that sorts the list of partition pairs (X, Y) by gain. For example, the nodes of partition X with partition Y as the candidate partition are ranked based on the gain that each node would have if it were remapped to the candidate partition. In this way, the node with the highest gain is first selected as the remapping target.
0069Considering Node Weights In one embodiment, the calculation of the total weights of the nodes in any two lists of partition pairs is to determine if the nodes can be remapped between these two partitions. You may. For example, if the total calculated weight of a node in partition # 3 is greater than the total weight of a node in partition # 4, and the node is in the list of partition pairs (3,4), then that partition pair (3,4) The nodes in the list may be remapped to partition # 4. The total weight of the nodes in the two partitions may be recalculated based on the updated remapping between the nodes in partitions # 3 and # 4, and the analysis may be repeated. If the total calculated weight of the node in partition # 3 is less than or equal to the total weight of the node in partition # 4, and the node exists in the list of partition pairs (4,3), the partition pair (4,3) The nodes in the list in 3) may be remapped to partition # 3. The total weights of the nodes in the two partitions may be recalculated based on the updated remapping and the analysis may be repeated. The comparison of total node weights for partitions and associated remappings when the node weights are not equal may be repeated any number of iterations. If the total node weights in each list of partition pairs are equal, or the differences between them are within the selection thresholds, then the analysis may be terminated to consider the partition to be load balanced. Good. In one embodiment, the total weight of the nodes in a partition is not only the updated remapping between the nodes in partitions # 3 and # 4, but also between the other partition and one or both of partitions # 3 and # 4. It may be recalculated based on any other updated remapping of the node. In general, the principles described herein for load balancing by node remapping are applicable to any combination of all partitions associated with a node in a node graph.
0070In one embodiment where the distributed system optimizes the mapping, the master computer system 313 generates a node list and a partition-to-partition (X, Y) matrix to determine how many nodes should be remapped to different partitions. You may try to do it. For example, master computer system 313 may receive a list of all partition pairs (X, Y) from aggregates 304-312. These lists may identify the total number of nodes in partition X with partition Y as a candidate partition. When considering node weights, the list may identify the total weights of the nodes in partition X with partition Y as the candidate partition.
0071The master computer system 313 may then maintain equilibrium by determining the number of nodes that should be remapped to different partitions. The master computer system 313 may then propagate the number of remapped nodes to computer systems 301-303. In one embodiment, to maintain equilibrium, if the list of two partition pairs (eg, partition pair (1,2) and partition pair (2,1)) has a different number of nodes, re-from each partition. The number of nodes for which mapping is selected is equal to the number of nodes in the list of partition pairs with the smaller number of nodes.
0072In one embodiment, the number of nodes remapped to different partitions is transmitted from master computer system 313 to computer systems 301-303 based on probability. For example, the probability of remapping a node in a list of partition pairs (X, Y) P<sub>XY</sub>May be determined based on the following equation. P<sub>XY</sub>= Z<sub>XY</sub>/ m<sub>XY</sub>Where Z<sub>XY</sub>Represents the number of nodes that are remapped from partition X to partition Y m<sub>XY</sub>Represents the total number of nodes in partition X with partition Y as a candidate partition.
0073Each node of partition X with partition Y as a candidate partition has a P probability of remapping to partition Y.<sub>XY</sub>Is. Termination module 403 may be calculated or analyzed to determine when the optimization process should be stopped. In one embodiment, the termination module 403 may stop the optimization process after the optimization has stabilized. In one embodiment in which the distributed system optimizes the mapping, the termination module 403 may be implemented by the master computer system 313. Further, the termination module 403 may determine whether or not the optimization process can be converged and stopped by executing the convergence detection technique. In one embodiment, the termination module 403 may stop the optimization process after a predetermined number of iterations.
0074In one embodiment, the termination module 403 is a "locality percentage" that can identify the total number or total weight of local edges (or edges within a partition to the entire partition) relative to the total number of edges in the node graph. Based on this, the optimization process may be stopped. The total number of edges may be calculated, for example, at the beginning of the optimization process. Termination module 403 may be made to determine the locality ratio by receiving global statistics on the node graph and node graph optimization from aggregates 304-312. Each node of a partition may contribute to the determination of the locality ratio based on whether its edges are local (in the same partition as the node). The local total edge weight may be determined by aggregating the local edge weights by the aggregates 301 to 303.
0075The locality ratio may be based on other considerations. Also, the locality ratio may be applicable to the total weight of the edges as opposed to the total number of edges. Termination module 403 may base the locality ratio on the ratio of the partition's local total weight to the total node weight of the node graph. In one embodiment, the determination of the locality ratio may take into account both the total number of local edges of the partition and the total weight of the edges of the partition.
0076Once the locality ratio stabilizes after a predetermined number of iterations, the optimization process may be considered convergent or stoptable. In one embodiment, such stability may be reflected in the change in the value of the locality ratio within the selection threshold at the selected number of iterations.
0077FIG. 5 shows an exemplary optimization method 500 for remapping nodes according to one embodiment. As a matter of course, the above description with respect to FIGS. 1 to 4 is also applicable to the method of FIG. Here, for the sake of brevity and clarity, we do not repeat any features and functions applicable to FIG.
0078In block 502 of method 500, the histogram of the nodes of the first partition may be calculated. In block 504, the histogram of the node of the second partition may be calculated. In one embodiment, blocks 502 and 504 may be implemented by the histogram generation module 201 of FIG. In one embodiment where the distributed system optimizes the mapping, each node produces a histogram based on the information transmitted from the machine associated with the connecting node.
0079In block 506, the second partition may be selected as a candidate partition for a set of nodes in the first partition based on the histogram of the nodes in the first partition. In block 508, the first partition may be selected as a candidate partition for a set of nodes in the second partition based on the histogram of the nodes in the second partition. In one embodiment, blocks 506 and 508 may be implemented by the candidate partition selection module 203 of FIG. Candidate partitions may be selected based on improved edge locality. In one embodiment, candidate partitions may be selected based on the probability of gain of edge locality. For example, candidate partitions may be selected based on a probabilistic algorithm that has a bias towards partitions with greater edge locality. In one embodiment, a local gain threshold may be implemented that must meet or exceed the threshold to qualify the partition as a candidate partition. The selection of candidate partitions for a node in a partition may involve one, many, or all partitions associated with a node in the node graph.
0080In block 510, based on the load balance of the partition, at least a part of the set of nodes in the first partition is remapped to the second partition and at least one of the set of nodes in the second partition. The part may be remapped to the first partition. In addition, load balancing by node remapping is applicable to some or all combinations of all partitions associated with nodes in the node graph. In one embodiment, block 510 may be implemented by the balanced remapping module 102 of FIG. In one embodiment that does not consider node weights, the total number of nodes remapped from one partition may be the same as the total number of nodes remapped from the other partition. Considering Node Weights In one embodiment, the total weight of the nodes remapped from one partition may be the same as the total weight of the nodes remapped from the other partition. In one embodiment where the non-distributed system optimizes the mapping, the nodes are remapped based on the gain so that the node associated with the highest gain is first selected for remapping. You may. In one embodiment, node remapping may be based on node weights. In one embodiment where a distributed system optimizes mapping, after the master computer system determines the number of nodes that should be remapped to each partition to maintain balance, the partition is present and remapped. The number may be communicated to the running computer system.
0081Social Networking System Illustrative Embodiment FIG. 6 is a network diagram of an exemplary system 600 that replaces a video link in a social network according to an embodiment of the present invention. The system 600 comprises one or more user devices 610, one or more external systems 620, a social networking system 630, and a network 650. In one embodiment, the social networking system described in connection with the above embodiments may be implemented as a social networking system 630. For convenience of description, an embodiment of system 600 shown in FIG. 6 includes a single external system 620 and a single user device 610. However, in other embodiments, the system 600 may include more user equipment 610 and / or one of more external systems 620. In certain embodiments, the social networking system 630 is operated by a social network provider, whereas the external system 620 is separated from the social networking system 630 in that it can be operated by different entities. ing. However, in various embodiments, the social networking system 630 and the external system 620 operate in concert to provide social networking services to users (members) of the social networking system 630. In this sense, the social networking system 630 provides a platform or backbone that can be utilized by other systems such as the external system 620 to provide users with social networking services and features throughout the Internet. Good.
0082The user device 610 includes one or more computer devices capable of accepting input from the user and transmitting / receiving data via the network 650. In one embodiment, the user device 610 is a conventional computer system running, for example, a Microsoft Windows compatible operating system (OS), Apple OS X, and / or a Linux distribution. In another embodiment, the user device 610 includes a smartphone, a tablet, and a personal digital assistant (PDA). Devices with computer functions such as Assistant) and mobile phones are possible. The user device 610 is configured to communicate over the network 650. In addition, the user device 610 can execute an application such as a browser application that enables the user of the user device 610 to interact with the social networking system 630. In another embodiment, the user device 610 interacts with the social networking system 630 through an application programming interface (API) provided by the user device 610's native operating system, such as iOS and ANDROID. User equipment 610 is external via network 650, which may include any combination of local area network and / or wide area network by using both wired and / or wireless communication systems. It is configured to communicate with system 620 and social networking system 630.
0083In one embodiment, the network 650 uses standard communication techniques and protocols. For this reason, network 650 includes Ethernet, 802.11, WiMAX (Worldwide Interoperability for Microwave Access), 3G, 4G, CDMA (Code Division Multiple Access), GSM (Global System for Mobile Communications), LTE (Long Term Evolution), and digital. -May include links using technologies such as Subscriber Line (DSL). Similarly, the networking protocols used in the Network 650 include Multi Protocol Label Switching (MPLS), Transmission Control Protocol / Internet Protocol (TCP / IP), and User Datagram. Protocol (UDP), Hypertext Transfer Protocol (HTTP: Hypertext) Examples include Transport Protocol (), Simple Mail Transfer Protocol (SMTP), File Transfer Protocol (FTP), and the like. Data exchanged over Network 650 shall be represented using technologies and formats, including Hypertext Markup Language (HTML) and Extensible Markup Language (XML), or technologies or formats. Can be done. In addition, all or part of the links are traditional, such as Secure Sockets Layer (SSL), Transport Layer Security (TLS), and Internet Protocol security (Ipsec). It can be encrypted using encryption technology.
0084In one embodiment, the user device 610 is a markup language document received from an external system 620 and a social networking system 630. By processing document) 614 with the browser application 612, content from the external system 620 and the social networking system 630 or the external system 620 or the social networking system 630 may be displayed. Markup language document 614 identifies the content and one or more instructions that describe the formatting or appearance of the content. By executing the instructions contained in the markup language document 614, the browser application 612 displays the identified content in the format or appearance described by the markup language document 614. For example, the markup language document 614 includes instructions to generate and display text and image data read from external systems 620 and social networking system 630 or a web page with multiple frames containing text or image data. .. In various embodiments, the markup language document 614 contains data files containing markup language data such as extended markup language (XML) data, extended hypertext markup language (XHTML) data, and the like. Be prepared. In addition, the markup language document 614 is intended to facilitate data exchange between the external system 620 and the user device 610 by including JSON (JavaScript Object Notation) data, JSONP (JSON with padding), and JavaScript data. You may do it. The browser application 612 on the user device 610 may use a JavaScript compiler to decrypt the markup language document 614.
0085The markup language document 614 may also include, and is linked to, applications or application frameworks such as FLASH or Unity applications, SilverLight application frameworks, and the like. May be good.
0086Also, in one embodiment, the user device 610 includes one or more cookies 616 containing data indicating whether or not the user of the user device 610 has logged in to the social networking system 630. It may be possible to change the data transmitted from the system 630 to the user device 610.
0087The external system 620 comprises one or more web servers including one or more web pages 622a, 622b that are propagated to the user equipment 610 using the network 650. Also, the external system 620 is separated from the social networking system 630. For example, the external system 620 is associated with the first domain, while the social networking system 630 is associated with a separate social networking domain. Web pages 622a, 622b included in the external system 620 include a markup language document 614 containing instructions that identify the content and specify the formatting or appearance of the identified content.
0088The social networking system 630 includes one or more users for the social network and provides the users of the social network with the ability to communicate and interact with other users of the social network. Equipped with computer equipment. In some cases, graphs, or data structures containing edges and nodes, can represent social networks. Social networks can also be represented using, but are not limited to, other data structures such as databases, objects, classes, meta elements, files, or any other data structure. The social networking system 630 may be operated, managed, or controlled by an operator. The operator of the social networking system 630 may be a human being, an automated application, or a set of applications that manage content, coordinate policies, and collect aggregate usage within the social networking system 630. The operator may be adapted to use any type.
0089After joining the social networking system 630, the user may add any number of connections with other users that the social networking system 630 wants to connect with. As used herein, the term "friend" refers to any other user of the social networking system 630 to whom a user has connected, relevance, or formed a relationship through the social networking system 630. For example, in one embodiment, when a user of the social networking system 630 is represented as a node in the social graph, the term "friend" is formed between the two user nodes and connects the two user nodes directly. Can represent an edge.
0090Connections may be added explicitly by the user or created automatically by the social networking system 630, based on the common characteristics of the user (eg, alumni of the same institution). Good. For example, the first user specifically selects a specific other user to be a friend. The connections in the social networking system 630 are usually bidirectional, but this is not required, so the terms "user" and "friend" are defined by the framework. The connections between users of the social networking system 630 are typically bidirectional (two-way) or mutual, but may be unidirectional or one-way. For example, if Bob and Joe are both users of social networking system 630 and are connected to each other, then Bob and Joe are connected to each other. On the other hand, if Bob wants to connect with Joe and see the data that Joe transmitted to social networking system 630, but Joe doesn't want to form a connection, a unidirectional connection is now established. You may be. The connection between users may be a direct connection. However, according to some embodiments of the social networking system 630, connections can be made indirectly by one or more connection levels or degrees of separation.
0091In addition to establishing and maintaining connections between users and allowing interaction between users, the social networking system 630 allows users to take action on various types of items supported by the social networking system 630. To do. These items include groups or networks to which users of social networking system 630 may belong (ie, social networks of people, entities, and concepts), events or calendar entries that users may be interested in, users. Computer-based applications available through the social networking system 630, services provided by the social networking system 630 or transactions that allow the user to buy or sell items through the social networking system 630, and the user social networking system. It may include interactions with ads that can run on or outside the social networking system 630. These are just a few examples of what users can do on the social networking system 630, and many others are possible. Users interact with anything that can be represented by social networking system 630 or external system 620, anything that is separate from social networking system 630, or anything that is coupled to social networking system 630 through network 650. It is possible to act.
0092In addition, the social networking system 630 can link various entities. For example, social networking system 630 allows users to interact with each other and with external systems 620 or other entities through APIs, web services, or other communication channels. The social networking system 630 generates and holds a "social graph" containing multiple nodes connected to each other by multiple edges. Each node in the social graph may represent an entity that can act on another node and / or an entity that can act on another node. The social graph may contain different types of nodes. Examples of node types include users, non-human entities, content items, web pages, groups, activities, messages, concepts, and any other thing that can be represented by objects in the social networking system 630. The edge between two nodes in the social graph may represent a particular kind of connection or association between two nodes that may result from a node relationship or an action taken by one node on the other node. Optionally, the edges between the nodes can be weighted. Edge weights can represent attributes associated with an edge, such as the strength of connections between nodes or relationships. Different types of edges allow different weights. For example, an edge created when one user likes another user is given some weight, while an edge created when a user becomes friends with another user is different. Weights may be given.
0093As an example, if the first user identifies the second user as a friend, the edge of the social graph connecting the node representing the first user and the second node representing the second user is generated. When the various nodes are related or interact with each other, the social networking system 630 modifies the edges connecting the various nodes to reflect the above relationships and interactions.
0094The social networking system 630 also includes user-generated content that enhances user interaction with the social networking system 630. User-generated content can include anything that users can add, upload, send, or "post" to social networking system 630. For example, a user propagates a post from user device 610 to social networking system 630. The post may include textual information such as recent status, location information, images such as photographs, videos, links, music, or other similar data and media or data or data such as media. The content may also be added to the social networking system 630 by a third party. Content "items" are represented as objects in social networking system 630. In this way, users of the social networking system 630 are encouraged to communicate with each other by posting text and content items in different types of media through different communication channels. Such communication increases the interaction of users with each other and the frequency with which users interact with the social networking system 630.
0095The social networking system 630 includes a web server 632, an API request server 634, a user profile store 636, a connection store 638, an action logger 640, an activity log 642, an authorization server 644, and an optimization module 646. In one embodiment of the invention, the social networking system 630 may include additional, fewer, or different components for a variety of applications. Other components such as network interfaces, security mechanisms, load balancers, failover servers, management and network operations consoles, etc. are not shown to avoid obscuring the details of the system.
0096User Profile Store 636 provides biographies, demographic information, and other types of descriptive information such as work history, educational background, hobbies or preferences, locations, etc. that have been declared by the user or inferred by the social networking system 630. , Holds information about user accounts. This information is stored in the user profile store 636 so that each user is uniquely identified. The social networking system 630 also stores data describing one or more connections between different users in the connection store 638. The connection information may indicate users with similar or common work history, group membership, hobbies, or educational background. The social networking system 630 also includes user-defined connections between different users that allow the user to specify relationships with other users. For example, according to user-defined connections, a user can generate relationships with other users, such as friends, colleagues, partners, etc., that are similar to the user's real-life relationships. The user may make a selection from a predetermined type of connection, or may specify the type of connection of the user, if necessary. Connections with other nodes of the social networking system 630, such as non-human entities, buckets, cluster centers, images, interests, pages, external systems, concepts, etc., are also stored in the connection store 638.
0097The social networking system 630 holds data about objects with which users can interact. To hold this data, the user profile store 636 and the connection store 638 store instances of the corresponding types of objects held by the social networking system 630. Each object type has an information field suitable for storing information appropriate for the object type. For example, the user profile store 636 contains a data structure with fields suitable for describing the user's account and information about the user's account. When a new object of a particular type is created, the social networking system 630 initializes the new data structure of the corresponding type, assigns a unique object identifier, and adds data to the object as needed. To start. This means that, for example, the user becomes a user of social networking system 630, social networking system 630 creates a new instance of the user profile in the user profile store 636, assigns a unique identifier to the user account, and the user provides it. This can happen if you start entering information in the fields of your user account.
0098The connection store 638 contains data structures suitable for describing a user's connection with other users, connection with an external system 620, or connection with other entities. The connection store 638 may also associate the type of connection with the user's connections that are available in conjunction with the user's privacy settings that coordinate access to information about the user. In one embodiment of the invention, the user profile store 636 and the connection store 638 may be implemented as a federated database.
0099According to the data stored in the connection store 638, the user profile store 636, and the activity log 642, the social networking system 630 uses nodes to identify various objects and uses the edges that connect the nodes. You can generate a social graph that identifies relationships between different objects. For example, in the social networking system 630, if the first user establishes a connection with the second user, the user accounts of the first and second users from the user profile store 636 will be in the social graph. It may act as a node. The connection between the first user and the second user stored in the connection store 638 is the edge between the first user and the nodes associated with the second user. Continuing with this example, the second user may send a message to the first user in the social networking system 630. The action of sending a storable message is another edge between the two nodes of the social graph representing the first user and the second user. The message itself may also be identified and included in the social graph as another node connected to the node representing the first user and the second user.
0100In another example, the first user tags the second user in an image held by the social networking system 630 (or an image held by another system outside the social networking system 630). You may attach it. The image itself may be represented as a node in social networking system 630. This tagging action creates an edge between the first user and the second user, as well as an edge between each user and the image (which is also a node in the social graph). May be good. In yet another example, if the user confirms attendance at the event, then the user and the event are the nodes obtained from the user profile store 636, and the attendance at the event is the node readable from the activity log 642. The edge between. By generating and retaining the social graph, the social networking system 630 contains a wealth of socially relevant information, including data describing different types of objects and the interactions and connections between them. source).
0101The web server 632 links the social networking system 630 to one or more user devices 610 and one or more external systems 620 or one via network 650. In addition to web pages, the web server 632 also provides other web-related content such as Java, JavaScript, Flash, and XML. The web server 632 may include messaging functions such as a mail server that receives and routes messages between the social networking system 630 and one or more user devices 610. The message can be an instant message, a queued message (eg, email), a text and SMS message, or any other suitable messaging format.
0102According to API request server 634, one or more external systems 620 and user equipment 610 can call access information from social networking system 630 by calling one or more API functions. Also, according to the API request server 634, the external system 620 can send information to the social networking system 630 by calling the API. In one embodiment, the external system 620 sends an API request to the social networking system 630 over the network 650, and the API request server 634 receives the API request. The API request server 634 generates an appropriate response by calling the API associated with this API request and processing the request, and the API request server 634 transmits this to the external system 620 via the network 650. .. For example, in response to an API request, the API request server 634 collects data associated with the user, such as the connection of the user logged in to the external system 620, and transmits this collected data to the external system 620. In another embodiment, the user device 610 communicates with the social networking system 630 via API in the same manner as the external system 620.
0103The action logger 640 is from a web server 632 for user actions on and outside social networking system 630 or user actions on or outside social networking system 630. Communication can be received. The action logger 640 also inputs information about the user action into the activity log 642 to perform various actions taken by the user of the social networking system 630 inside and outside the social networking system 630. Allow system 630 to discover. Any action taken by a particular user with respect to another node on social networking system 630 is a data repository such as activity log 642 or a similar database. The information stored in the repository) can be associated with each user's account. Examples of user identification and storage actions in social networking system 630 include, for example, adding a connection to another user, sending a message to another user, reading a message from another user, and so on. This includes browsing content associated with another user, attending events posted by another user, posting images, attempting to post images, or other actions that interact with another user or another object. When a user takes an action on the social networking system 630, the action is recorded in the activity log 642. In one embodiment, the social networking system 630 maintains activity log 642 as a database of entries. When an action is taken in the social networking system 630, an entry for that action is added to the activity log 642. The activity log 642 is sometimes referred to as the action log.
0104User actions may also be associated with concepts and actions that occur within an entity external to the social networking system 630, such as an external system 620 separate from the social networking system 630. For example, the action logger 640 may receive data from the web server 632 that describes the user's interaction with the external system 620. In this example, the external system 620 reports user interactions according to the structured actions and objects of the social graph.
0105Other examples of actions in which the user interacts with the external system 620 are that the user indicates interest in the external system 620 or another entity, and that the user socializes comments about web page 622a within the external system 620 or external system 620. Posting to Networking System 630, User Posting Uniform Resource Locator (URL) or other identifier associated with External System 620 to Social Networking System 630, Event Associated with External System 620 The user may be present, or any other action by the user regarding the external system 620. Thus, the activity log 642 may include an action that describes the interaction between the user of the social networking system 630 and an external system 620 that is separate from the social networking system 630.
0106The authorization server 644 enforces one or more privacy settings for users of social networking system 630. The user's privacy settings determine how certain information associated with the user can be shared. Privacy settings include specifications for specific information associated with users and specifications for one or more entities with which information can be shared. Examples of entities that can share information include other users, applications, external systems 620, or any entity that has potential access to information. Information that can be shared by the user includes user account information such as profile pictures, telephone numbers associated with the user, user connections, addition of connections, changes to user profile information, and other actions taken by the user.
0107Privacy settings specifications have different levels of It may be provided in granularity). For example, privacy settings may identify certain information that is shared with other users, and certain relevant information such as work phone numbers or profile pictures, home phone numbers, and personal information, including status. Identify the set. Alternatively, the privacy settings can be applied to all information associated with the user. The specifications of entity sets that can access specific information can also be specified at various particle size levels. Various entities that can share information include, for example, all friends of a user, all friends of friends, all applications, or all external systems 620. According to one embodiment, the entity set specification can include an entity list. For example, the user may provide a list of external systems 620 that are allowed access to certain information. According to another embodiment, the specification can include a set of entities along with exceptions that do not allow access to the information. For example, the user may allow all external systems 620 to access the work information, but may specify a list of external systems 620 that are not allowed access to the work information. In certain embodiments, a list of exceptions that are not allowed access to certain information is referred to as a "block list". The external system 620 belonging to the block list specified by the user is blocked from accessing the information specified in the privacy setting. Various combinations are possible for the particle size of the information specification and the particle size of the specification of the entity with which the information is shared. For example, all personal information may be shared with friends, while all work information may be shared with friends of friends.
0108Authorization server 644 includes logic that determines whether a user's friends, external system 620, and / or other applications and entities have access to specific information associated with the user. The external system 620 may require approval from the approval server 644 to access the user's more private and sensitive information, such as the user's work phone number. Based on the user's privacy settings, the authorization server 644 has access to information associated with the user, such as information about actions taken by the user, by another user, external system 620, application, or another entity. Judge whether or not.
0109The optimization module 646 may be made to optimize the current mapping of nodes in the social graph with respect to the social networking system 630. The optimization module 646 may also optimize the current mapping to improve edge locality while maintaining load equilibrium. In one embodiment, the optimization module 646 may be implemented as the optimization module 100 of FIG. Hardware Embodiment The above processes and features can be implemented in a wide variety of network and computer environments with a wide variety of machine and computer system architectures. FIG. 7 shows an example of a computer system 700 that can be used in one or more implementations of the embodiments described herein according to an embodiment of the present invention. The computer system 700 includes a set of instructions that causes the computer system 700 to perform the processes and features described herein. The computer system 700 may also be connected to another machine (eg, networked). In a network deployment, the computer system 700 is a server machine or client machine in a client / server network environment or a peer machine (or peer) in a peer-to-peer (or distributed) network environment. It may operate as a machine). In one embodiment of the invention, the computer system 700 may be a component of the social networking system described herein. In one embodiment of the invention, the computer system 700 may be one of many servers that make up all or part of the social networking system 630.
0110The computer system 700 is housed in a processor 702, a cache 704, and a computer-readable medium and comprises one or more executable modules and drivers for the processes and features described herein. The computer system 700 also includes a high performance input / output (I / O) bus 706 and a standard I / O bus 708. The host bridge 710 couples the processor 702 to the high performance I / O bus 706, while the I / O bus bridge 712 couples the two buses 706 and 708 to each other. The high-performance I / O bus 706 is coupled with system memory 714 and one or more network interfaces 716. The computer system 700 may further include a video memory and a display device coupled to the video memory (not shown). Mass storage 718 and I / O port 720 are combined with standard I / O bus 708. The computer system 700 may optionally include input / output devices (not shown) such as a keyboard and pointing device, display device, etc. coupled to the standard I / O bus 708. These elements are generally x86 compatible processors manufactured by Intel Corporation in Santa Clara, California, Advanced Micro Devices (AMD), Inc. in Sunnyvale, California. ) Is intended to represent a wide range of computer hardware systems, such as, but not limited to, x86 compatible processors manufactured by) and computer systems based on any other suitable processor.
0111The operating system manages and controls the operation of the computer system 700, such as inputting and outputting data to software applications (not shown). The operating system provides an interface between software applications running on the system and the hardware components of the system. LINUX operating system, Apple Macintosh operating system, UNIX operating system, Microsoft® Windows® operating system available from Apple Computer Inc., Cupertino, Calif. Any suitable operating system, such as the BSD operating system, can be used. In addition, other embodiments are also possible.
0112The elements of the computer system 700 will be described in more detail below. In particular, network interface 716 provides communication between the computer system 700 and any of a wide range of networks, such as Ethernet (eg, IEEE 802.3) networks, backplanes, and the like. Mass storage 718 provides storage of data and programming instructions that perform the processes and features implemented by each of the identified computing systems. System memory 714 (eg, DRAM), on the other hand, provides temporary storage of data and programming instructions when executed by processor 702. I / O port 720 is one or more serial and parallel communication ports or one or more serial communication ports or parallel communication that provide communication between additional peripherals that can be coupled to computer system 700. It may be a port.
0113The computer system 700 may include various system architectures, and various components of the computer system 700 may be rearranged. For example, the cache 704 may be on-chip with the processor 702. Alternatively, the cache 704 and the processor 702 may be integrally packaged as a "processor module", the processor 702 being referred to as a "processor core". Furthermore, in a particular embodiment of the present invention, all of the above components may or may not be required. For example, the high performance I / O bus 706 may be coupled with peripherals coupled to the standard I / O bus 708. Also, in some embodiments, there may be only one bus and the components of the computer system 700 may be combined into this single bus. In addition, the computer system 700 may include additional components such as additional processors, storage devices, or memory.
0114In general, the processes and features described herein may be implemented as part of an operating system or a set of instructions referred to as a particular application, component, program, object, module, or "program." .. For example, the particular process described herein may be performed using one or more programs. A program typically gives a computer one or more instructions to cause the computer system 700 to perform operations that perform the processes and features described herein when read and executed by one or more processors. · Included in various memory and storage devices in System 700. The processes and features described herein may be implemented in software, firmware, hardware (eg, application-specific integrated circuits), or any combination thereof.
0115In one embodiment, the processes and features described herein may be implemented individually or collectively as a series of executable modules operated by the computer system 700 in a distributed computer environment. The module may be implemented by hardware, an executable module stored on a computer-readable medium (or machine-readable medium), or a combination thereof. For example, these modules may include multiple instructions or a set of instructions executed by a processor in the hardware system, such as processor 702. Initially, a series of instructions may be stored in a storage device such as mass storage 718. However, the series of instructions can be stored in any suitable computer-readable storage medium. Further, the series of instructions need not be stored locally and can be received from a remote storage device such as a server on the network via the network interface 716. The instruction is accessed and executed by the processor 702 after copying from a storage device such as mass storage 718 to system memory 714. In various embodiments, one or more modules can be executed by one or more processors in one or more locations, such as multiple servers in a parallel processing environment.
0116Examples of computer-readable media include recordable types of media such as volatile and non-volatile memory devices, solid state memory, removable disks such as floppy disks, hard disk drives, magnetic media, optical disks (eg, compact disc reads). By execution by a single memory (CD ROM), digital versatile disc (DVD)), or other similar non-volatile (or temporary) tangible (or intangible) storage medium, or computer system 700. Any kind of medium suitable for storing, encoding, or transmitting a series of instructions performing any one or more of the processes and features described in, but is not limited thereto.
0117For convenience of explanation, many specific details have been provided to provide a thorough understanding of the present specification. However, it will be apparent to those skilled in the art that the embodiments of the present disclosure are feasible without these specific details. In some cases, modules, structures, processes, features, and equipment are shown in the form of block diagrams to keep the description from being confusing. In another example, the flow of data and logic is represented by showing a functional block diagram and a flow diagram. The components of block diagrams and flow diagrams (eg, modules, blocks, structures, equipment, features, etc.) are various combinations, separations, removals, rearrangements, in ways not expressly described and illustrated herein. And may have been replaced.
0118References in the present specification such as "one embodiment", "another embodiment", "a series of embodiments", "some embodiments", "various embodiments" and the like will be described in relation to the embodiments. It is meant that the particular feature, design, structure, or property is included in at least one embodiment of the present disclosure. The appearance of expressions such as "in one embodiment" in various parts of the specification does not necessarily represent the same embodiment, nor is it a separate or separate embodiment that mutually excludes other embodiments. Furthermore, various features that can be included in various combinations in some embodiments and can be omitted in various ways in other embodiments, regardless of whether or not there is an explicit reference such as "embodiment". Is described. Similarly, it describes various features that may be preferences or requirements in some embodiments but not in other embodiments.
0119The expressions used herein are selected primarily for readability and teaching convenience, and may not be selected to describe or limit the subject matter of the present invention. Therefore, the scope of the present invention shall be limited not by this detailed description but by any claim derived from the use based thereto. From the above, the disclosure regarding the embodiment of the present invention is an example of the scope of the present invention shown in the following claims, and is not limited in any way.
7 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7
Every citation, both ways
| Document | Relation | Office |
|---|---|---|
| US20060015588A1 | Cites | United States of America |
| WO2011037505A1 | Cites | World Intellectual Property Organization (WIPO) |
| JP2013126158A | Cites | Japan |
19 members in 11 offices
Members19
| Document | Office | Kind | |
|---|---|---|---|
| US2015095348A1 | United States of America | A1 | |
| EP2858026A2 | European Patent Office (EPO) | A2 | |
| CA2925114A1 | Canada | A1 | |
| WO2015050568A1 | World Intellectual Property Organization (WIPO) | A1 | |
| EP2858026A3 | European Patent Office (EPO) | A3 | |
| AU2013402190A1 | Australia | A1 | |
| IL244779A0 | Israel | A0 | |
| CN105593838A | China | A | |
| KR20160062073A | Republic of Korea | A | |
| MX2016004118A | Mexico | A | |
| JP2016540280A | Japan | A | |
| BR112016007205A2 | Brazil | A2 | |
| JP6272467B2This record | Japan | B2 | |
| US9934323B2 | United States of America | B2 | |
| MX355951B | Mexico | B | |
| CN105593838B | China | B | |
| CA2925114C | Canada | C | |
| KR102076580B1 | Republic of Korea | B1 | |
| AU2013402190B2 | Australia | B2 |
13 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Cancellation because of no payment of annual feesLAPS | LAPS | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Written notification of registration of transferJAPANESE INTERMEDIATE CODE: R350R350 | R350 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Written request for registration of change of nameJAPANESE INTERMEDIATE CODE: R313533S533 | S533 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Notification of acceptance of power of attorneyJAPANESE INTERMEDIATE CODE: R3D02RD02 | RD02 | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| First payment of annual fees (during grant procedure)JAPANESE INTERMEDIATE CODE: A61A61 | A61 | |
| Written decision to grant a patent or to grant a registration (utility model)JAPANESE INTERMEDIATE CODE: A01A01 | A01 | |
| Decision of grant or rejection writtenTRDD | TRDD | |
| Report on retrievalJAPANESE INTERMEDIATE CODE: A971007A977 | A977 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 6272467
- Application
- 2016520081
Titles2
- Japanese
- 局所性および平衡のための動的マッピングのシステムおよび方法
- English
- Dynamic mapping systems and methods for locality and equilibrium
Classification
- CPC, 3
- G06F16/9024
- G06Q10/101
- G06Q10/48
- IPC, 1
- G06F17 30
