Distributed cache for graph data
Abstract
A distributed caching system for storing and serving information modeled as a graph that includes nodes and edges that define associations or relationships between nodes that the edges connect in the graph.

Term
Projected expiry 30 November 2031.
- Priority
- Filed
- Granted
- Today
- Projected expiry
18 claims: 3 independent, 15 dependent
- 1Claims 1, A system comprising:one or more first computing devices providing a persistent-storage database operative to maintain a graph data structure comprising a plurality of graph nodes and a plurality of graph edges connecting the graph nodes, a graph edge connecting two graph nodes indicating an association between the two graph nodes, each graph node being a data object corresponding to a profile associated with a social-networking system and having a unique graph-node identifier;and a plurality of second computing devices coupled to the one or more first computing devices and providing a cache layer between the persistent-storage database and a plurality of client servers, the cache layer comprising a plurality of follower cache clusters that each comprise one or more follower cache nodes, each follower cache node comprising one or more one individual computing system, each follower cache node being operative to: maintain in the follower cache node at least a portion of the graph data structure, wherein the portion of the graph data structure comprises a plurality of graph nodes and □ plurality of graph edges, and wherein a count value of an association set is maintained for each graph node, the count value being either increased or decreased in response to commands to add or delete an association with respect to the graph node, respectively;receive a query from a user of the social-networking system for associations between nodes in the portion of the graph data structure maintained in the follower cache node, wherein the user is associated with a particular profile that corresponds to a graph node in the portion of the graph data structure maintained in the follower cache node;and respond to the query for associations between nodes in the graph data structure at least in part by accessing the portion of the graph data structure maintained in the follower cache node.
- 68. A method comprising:by one or more first computing devices, providing a persistent-storage database operative to maintain a graph data structure comprising a plurality of graph nodes and a plurality of graph edges connecting the graph nodes, a graph edge connecting two graph nodes indicating an association between the two graph nodes, each graph node being a data object corresponding to a profile associated with a social-networking system and having a unique graph-node identifier;and by a plurality of second computing devices coupled to the one or more first computing devices and providing a cache layer between the persistent-storage database and a plurality of client servers, the cache layer comprising a plurality of follower cache clusters that each comprise one or more follower cache nodes, each follower cache node comprising one or more one individual computing system, each follower cache node being operative to: maintaining in the follower cache node at least a portion of the graph data structure, wherein the portion of the graph data structure comprises a plurality of graph nodes and a plurality of graph edges, and wherein a count value of an association set is maintained for each graph node, the count value being either increased or decreased in response to commands to add or delete an association with respect to the graph node, respectively;*11540696 V2 CA 2974065 2019-01-22 receiving a query from a user the social-networking system for associations between nodes in the portion of the graph data structure maintained in the follower cache node, wherein the user is associated with a particular profile that corresponds to a graph node in the portion of the graph data structure maintained in the follower cache node;and responding to the query' for associations between nodes in the graph data structure at least in part by accessing the portion of the graph data structure maintained in the follower cache node.
- 1315. Λ plurality of non-transitory computer-readable storage media embodying software that is operative when executed to:provide a persistent-storage database operative to maintain a graph data structure comprising a plurality of graph nodes and a plurality of graph edges connecting the graph nodes, a graph edge connecting two graph nodes indicating an association between the two graph nodes, each graph node being a data object corresponding to a profile associated with a socialnetworking system and having a unique graph-node identifier;and CA 2974065 2019-01-22 provide a cache layer between the persistent-storage database and a plurality of client servers, the cache layer comprising a plurality of follower cache clusters that each comprise one or more follower cache nodes, each follower cache node comprising one or more one individual computing system, each follower cache node being operative to;maintain in the follower cache node at least a portion of the graph data structure, wherein the portion of the graph data structure comprises a plurality of graph nodes and a plurality of graph edges, and wherein a count value of an association set is maintained for each graph node, the count value being either increased or decreased in response to commands to add or delete an association with respect to the graph node, respectively;receive a query from a user of the social-networking system for associations between nodes in the portion of the graph data structure maintained in the follower cache node, wherein the user is associated with a particular profile that corresponds to a graph node in the portion of the graph data structure maintained in the follower cache node;and respond to the query for associations between nodes in the graph data structure at least in part by accessing the portion of the graph data structure maintained in the follower cache node,
Independent claims3
128 paragraphs in 11 sections, as filed
WO 21)12/091846
PCT/US2011/062609
DISTRIBUTED CACHE FOR GRAPH DATA
TECHNICAL FIELD
The present dise leisure relates generally lo storing and serving graph daia, and more parti till a rly, tv storing and serving graph data with a distributed cache system.
BACKGROUND
Computer users are able to access and sliarre vest amounts of information through various local and wide area computer networks including proprietary networks as well as public networks such as the Internet. Typically, a web browser installed on n user's computing device facilitates access to and interaction wills information located at various network servers identified by, for example, associated uniform resource locators fURLs). Conventional approaches to enable shoring of uscr-gcncmlcd content include various information sharing technologies or phitfomis such as social networking websites. Such .websites may include, be linked with, or provide a platform for applications enabling users to view web pages created or customized by other users where visibility and interaction with such pages by Other Users is governed by sonic characteristic set of rules,
Such social networking information, and most information in general, is typically stored in relational databases. Generally, a relational database is a collection of relations ffrequently referred -U as tables). Relational dautbiises use » set of mathematical terms, which may use Structured Query Language (SQL) database terminology. For example, a relation may be defined as a set of tuples that have the same attributes. A tuple usually represents an object sud ΐπΓυτιιιηΙίιπι jihotil that object. A relation ts usually described as a table, which is organized into rows and columns. Generally, all Hie data relcienccd by an attribute are in the same domain and conform to tlx* same constraints.
The relational model specifies that the tuples of a relation have no specific order and that (he tuples, in him, impose no order on die attributes. Applications access data by specifying queries, which use operations to identify tuples, identify attributes, turd to combine relations. Relations can be modified and new tuples can supply explicit values or be derived from a query.
CA 2974065 ÏO17-O7-21
Siinilaïiy, queries identify may tuples fur updating or deleting. It is necessary fur each tuple uf a relation lu be uniquely identifiable by some combination (one or more) of its attribute values. This combination is referred to as the primary key. In a relational database, al! data are stored ant! accessed via relations. Relations that store data arc typically implemented with or referred to us tables. .
Relational databases, as implemented in relational database management systems, have become a predominant choice for the storage of information in databases used for, for example, financial records, manufacturing and logistical information, personnel data, and other applications. As computer power has increased, die inefficiencies of relational databases. which
If) made them impractical in earlier times, have been outweighed by their case of use fur conventional applications. The three leading open source implementations are MySQL. I'ostgieSQL, and SQLife. MySQL is a relational database management system (RDBMS) that runs as a server providing multi-user access to a number of databases. The M” tn the acronym of the popular LAMP software stack refers to MySQL. Its popularity for tree with web
Ii applications is closely tied to the popularity of PI-IP (the P in I.AMP). Several high-traffic web sites use MySQL for data storage and logging of user data.
As communicating with relational databases is often a speed bottleneck, many networks ulilive caching systems to serve particular information queries. For example. Mcmesched is a general-puiposc distributed memory caching system. It is often used lo speed
2tJ up dynamic database-driven websites by caching data and objects in RAM to reduce tire number of times an external tinta source (such as a database or API) must be read. Memcachcd's APIs provide a giant bash table distributed across multiple machines. When die table is full, subsequent inserts cause older data to be purged in least recently treed (LRU) order. Applications Using Memcached typically layer requests and additions into core before falling back on a slower backing store, such as a database.
Tbe Memeaehed system uses a client server architecture. The seiveiv maintain u key value associative nrruy; the clients populate this array and query it. Clients use client side libraries lo contact the servers. Typically, each client knows all servers and the servers du not communicate with each other. If a client wishes to set or read die value cunespuudiiig to ;i txnniit key, the client's library first computes a hash of the key lo determine the server dint will
CA 2S14OG5 2017-07-21 tx· used. the client then contacts that server. The server will compute a second hash of the key to determine where io store or read the corresponding value. Typically, rite servers keep the values in RAM; if a server inns out of RAM, it discards the oldest vu lues. Therefore, clients must treat Memcaclicfl as a transitory cache; they cannot assume that data stored in Mcmcached is still there when they need it.
BRIEF DESCRIPTION OF THE DRAWINGS
Figure I illustrates an example caching system architecture according to one implementation of the invention.
Ill Figure 2 illustrates an example computer system architecture.
Figure 3 provides an example network environment.
Figure 4 shows a floweitart illustrating an example method (or adding a new association to a graph.
Figure S is a schcmalic diagram illustrating an example message flow between various components of a caching system.
Figure ή shows a flowchart illustrating an example method For processing changes to graph data.
Figure 7 is a schematic diagram illustrating an example message flow between various eoiitpoiteuls oFa caching System.
DESCRIPTION OF EXAMPLE EMBODIMENTS
Particular embodiments relate to a distributed caching system for storing and serving information modeled us a graph (hat includes nodes and edges that define associations or relationships between nodes that the edges connect in the graph, hi particular embodiments, the graph is. or includes, a sociai graph, and lhc distributed caching system is part of a huger networking system, infrastructure, or platform that enables an integrated social network environment. In the present disclosure, the social network environment may be described in terms of □ social graph including sociai graph information. In fact, particular embodiments of the present disclosure rely on, exploit, or make use of the fact thal most or al) of the tlata stored by or for the social network environment can be represented as n social graph. Particular
CA 2974065 2017-07-21 embodiments provide a cost-effective infrastructure that can efficiently, intelligently, and successfully scale with the exponentially increasing uttinbcr of users of il'c social network environment stteh as that described herein.
In particular enihndimenls, (he distributed caching system and backend infrastructure described herein provides one or more of: low latency at scale, a lower cost per request, an easy to use framework for developers, an infrastructure that supports multi-inns ter, un infrastructure (hat provides access to stored data to clients written in languages other than Hype (text Preprocessor (PHP), an infrastructure (hat enables combined queries involving both associations (edges) and objects (nodes) of a social graph as described hy way of example herein, and an infrastructure that enables diflcreni persistent data stores to be used for different types of data. Furthermore, particular embodiments provide one or more of: an infrastructure that enables a clean sepitniliun of tine data access API from the cuching+persistcrce+repliealion infrastructure, an infrastructure that supports write-lliruugh/rcud-through caching, an infrastructure that moves computations closer to tire data, an infrastructure that enables transparent migration to different storage schemas and book ends, and an infrastructure that improves the efficiency of data object access,
Additionally, as used lrerein. “or” may imply “and as well as or; that is, or does not necessarily preclude “and. unless explicitly stated or implicitly implied.
Panirular embodiments inay operate in a wide area network environment, such as die
2ft Internet, including multiple network addressable systems. Figure 3 illustrates an example network environment. tn which various example embodiments may operate. Network cloud 60 generally represents one or more interconnected networks, over which the systems and hosts described Ire re in can communicate. Network cloud 60 inay include packet-bused wide area networks (such as the Internet), private networks, wireless networks, satellite networks, cellular networks, paging networks. and the like. As Figure 3 illustrates, particular embodiments may Operate in a network environment comprising social networking system 20 and one or more client devices 30. Client devices 30 arc operably connected to the network environment via o network service pro vider, a wireless carrier, or any other suitable means.
In one example embodiment, social networking system 20 comprises computing sysiems dial allow users to communicate or otherwise idler a cl with each oilier and access
CA 2S74OG5 2017-07-21 content, such as user profiles, as described herein. Social networking system 20 is a network addressable system that, in vurious example embodiments. comprises one or mere physical servers 22 ant! data store 24. The one or more physical servers 22 arc operably connected to computer network 60 via. by way of example, a set of routers and/or networking switches 26. In an example embodiment, the functionality hosted by the one or more physical servers 22 may include web or HTTP servers, FTP servers, as well as, without limitation, web pages and applications imp lenient eti using Common Gateway Interface (CGI) script, PHP Hyper-text Preprocessor (PHP), Active Server Pages (ASP), Hyper Text Markup Language (HTML), Extensible Markup Language (XML). Java, JavaScript, Asynchronous JavaScript and XML (AJAX), and the like.
Physical set'vers 22 may host functionality directed lu tin.· operations uf social networking system 20. By way of example, so ci ill networking system 20 may I tost a website that allows oik; oi unne users, at one or inoie client devices 30, lo view and post information. hs well as ewminiiuientc with one another via the website, flercinafler servers 22 may be referred lo as server 22, although server 22 may include numerous servers hosting, for example, social networking system 20, as well ns other content distribution servers, data stores, and databases. Data store 24 may store content and data relating to. and enabling, operation of the social networking system as digital data objects. Λ data object, in particular implementations, is an iietn of digital information typically stored or embodied in a data file, database or record.
Content obtecis may take many forms, including: text (e.g., ASCII. SGML, HTML), images ie.W.. jpeg, lif and gif), graphics (vector-based or bitmap), audio, video (e.g., mpeg), or other multimedia. and combinations thereof. Gcuicnf object data may also include executable code objects (e.g., games executable within a browser window or frame). podcasts, etc. Logically, data store 24 corresponds to one or mote of a variety of separate and integrated databases, such as relational databases and object-oriented databases, that maintain information as an integrated collection of logically related records or files stored on one or more physical systems. Structurally, data store 24 may generally include one or more of n large class of data storage and management systems. In particular embodiments, dala store 24 may be implemented by tiny suitable physical system(s) including components, such as one or more database servers, mass .1(1 storage media, media library systems, storage area networks, data storage clouds, and lire like. Lu
CA 2374O65 ÎO17-O7-21 one example embodiment, dala store 24 includes one ot more servers, databases (e.g., MySQL), and/or data warehouses.
Data store 24 may include diihi associated with different social netwoiking system 20 users and/or client devices 30. In particular embodiments, the social networking system 20 nuinbuns a user profile Tor each user of the system 20. User profiles include data that describe ihe users of a social network, which may include, fur example, proper names (first, middle and last of a person, a trade name and/or company name of a business entity, etc.) biographic, demographic, and other types of descriptive information, such as work experience, educational history, hobbies or preferences, geographic location, and additional descriptive data. By way of example, user profiles may include a user’s birthday, relationship status, city of residence, and the like. The system 20 may further store data describing one or more relationships between different users. The relationship information may indicate users who have similar or common work experience. group memberships, hobbies, or cduenttonal history. Λ user profile may also include privacy settings governing access to Ihe user’s information is ro other users.
Client device 30 is generally a computer or computing device including functionality for communicating (e.g., remotely) over a computer network. Client device 30 may be a desktop computer, laptop computer, personal digital assistant (PDA), in- or out-of-ear navigation system, smart phone or other cellular or mobile phone, or mobile gaming device, among other suitable computing devices. Client device 30 may execute one or more client applications, sti eh ns a web browser (e.g., Microsoft Windows Internet Explorer, Mozilla Firefox. Apple Safari, Google Chrome, and Opera, etc,), to access and view eorttcnl over o eonipulcr network. In particular implementations. the clienl applications allow a user of client device 30 to enter addresses of specific network resources lo be retrieved, such as resources hosted by social networking system 20. These addresses can be Uniform Resource Locators, or UPJ-3. Irt addition, once a page ot oilier resource has been retrieved, Ihe client applications may provide access to oilier pages or records when the user ''clicks on hyperlinks to other resources. By way of example, such hyperlinks may be located within the web pages and provide an automated way for the user to enter the URL of anolher page and lo retrieve that page.
Figures I illustrates an example embodiment of a networking system, aiebiteeture, or miiasinieruie 100 (hereinafler referred to as networking syslcni 100) that enn implement the
CA 2Ï74OS5 2017-07-21 tack end functions of social networking system 20 illustrated in Figure 3. In particular embodiments, networking system 100 enables users of networking system 100 to interact with each other via social networking sendees provided by networking system 100 ns well as with third parties. For example, users at remote user computing devices (e.g.. personal computers, netbooks, multimedia devices, cellular phones (especially smart phones), etc.) may access ttchvorkinu system ItHl via web browsers or other user client applications to access websites, web pnges. or web applications hosted or accessible, al least in purt, by networking system I()0 to view information, store or update information, communicate information, or otherwise intemet with other users, third party websites, web pages, or web applications, or oilier information stored, hosted, or accessible by networking system 100. in pari ten lar embodiments, networking system 100 maintains a graph that includes graph nudes representing users, concepts, topics, anti other information (data), as well ns graph edges that connect or define relationships between graph nodes, as described in more detail below.
With reference lo Figures I arid 5, in particular embodiments, networking system (00 includes one or more data centers 102. Fur example, networking system 100 may include a plurality of data centers 102 located strategically within various geographic regions for serving users located within respective regions. In particular embodiments, each data eenter includes a number of client or web servers Î04 (hereinafter client servers 104) lhat communicate information to and from- users of networking system I(XI. For example, users at remote user computing devices may commun ic ate with client servers iO4 via load balancers or other suitable systems via any suitable combination of networks and service providers. Client servers 104 me y query tile caching system described herein in order to retrieve data to generate structured documents lor responding to user requcsis.
Each of the client servers 104 communicates wilh one or more follower distributed cache clusters or rings 106 (hereinafter Fol .ower cache clusters 10b). In the illustrated embodiment, data center 102 includes three follower cache clusters 106 that each serve a subset of the web server s 104. In particular embodiments, a follower cache cluster 106 and the client servers 104 lire follower cache cluster 106 serves are located inclose proximity, such as within a building, room, or other centralized location, which ted tees costs associated with the infrastructure (e.g., wires or other communication lines, etc.) as well ns latency between the
CA 2Ϊ74ΟΕ5 2017-07-21 client servers 104 and respective serving follower cache nodes cluster 106. However, in sonic embotiimentsl while each of the follower cache clusters 106, and the client servers 104 they respectively serve, may be located within a centralized location, each of the follower cache clusters 106 and respective client servers 104 tire follower cache clusters 106 respectively save, ί may be located in a different location than the other follower cache clusters 106 and respective client servers 104 of a given data center; that is. the follower cache clusters 106 (and the respective client servers 104 the clusters serve) of a given data center of a given region may be disiriborcd throughout various locations within the region.
hi particular embodiments, each data center 102 further includes a leader cache
III cluster I DR tliai communicates information between the follower cache clusters I Dû of a given data cellier i 02 and a persistent storage database I 10 oftltc given data center 102. In particular cmboditnellls, dal abuse 110 is a relational database. In particular embodiments, leader cache cluster IOS may include a plug-in operative to interoperate with any suitable implement a I ion of database 110. For example, database I ID may be implemented ns a dynamically-variablc plug-in
IS architecture and may utilize MySQL, and/or any suitable relational database management system such as, for example, HAYSTACK, CASSANDRA, among others. In one implementation, the plug-in performs various translation operations, such as translating data stored in tire caching layer as graph noties anti edges to queries anti commands suitable for n relational database including one or more tables or Hal files. In particular embodiments, leader cache cluster 108
2If also coordinates write requests lo database 110 from follower cache clusters 106 and sometimes read requests from follower eaohe clusters 106 for information cached in leader cache cluster 108 or (if not cached in leader cache cluster IOS) stored in database 110. In particulnr embodiments, leader cache cluster 108 furtlrcr coordinates the synchronization of information stored in the follower cache clusters 106 Of the respective da la center 102. That is. in particular embodiments.
the leader cache cluster IOS of a given dutii center 102 is con figurer I io maininin cache consistency (e.g.. the information cached) between the follower cache clusters 106 of the data center 102. to maintain eaehe consistency between the follower cache clusters 106 and the leader cuehv cluster I OK. and to slore the information cached in leader eaehe cluster 108 wilhin dutubuse H0. In one implementation, it leader cache cluster 108 and a follower cache cluster 106 can be
3Ü considered a caching layer between client servers 104 and database 1 10,
CA 2S74OG5 ÎO17-O7-21 lit one implementation, the caching layer is a write-tliru/rcatf*thru caching layer, wherein all reads and writes traverse the caching layer. In one implementation, die caching layer maintains association information and, thus, can handle queries for such information. Other queries ate passed through to database 110 for execution. Database 110 generally connotes a database system that may itself include other caching layers for handling other query types.
Each follower cache cluster 106 may include a plurality of follower cache nodes I 12, each of which may be running On an individual computer, computing system, or server. However, as described ahovc, each of the follower cache nodes 112 of a given follower cache cluster 10b may he located within π centralized location. Similarly, each leader cache cluster
If) 108 may include a plurality of leader cache nodes 114. each of which may be running on an individual computer, computing system, or server. Similar lo the follower cache nodes 112 of a given follower cache cluster 106, each of the leader cache nodes 114 of a given leader cache cluster 108 tnay be located within a ccnlraliz-cd location. For example, each dala eenter 102 may include tens, hundreds, or thousands of client servers 104 and each follower cache cluster 106
IS may include tens, hundreds, or thousands of follower cache nodes 112 that serve a subset of the them servers 104. Similarly, tach leader cache cluster IOS may include lens, hundreds, or thousands ol' leader cache nodes 114. In pnrlicuiiti embodiments, each of the follower cache nodes 112 within a given follower cache cluster 106 may only communicate wilh ihe other follower cache nodes 112 within .the particular follower cache cluster 106, tire client servers 104 served by Lite particular follower cache cluster 106, and the leader cache nodes 1 )4 within tbe leader cache cluster 108..
In particular embodiments. -,nforeunion stored by networking system 100 is stored within each data center 102 both within database 110 as well as within each ofthe follower and leader cache clusters 106 and 108, respectively, in particuiar embodiments, the infonnation stored within cadi database 110 is stored rehuiona I ly (e.g., as objects and tables via MySQL), uhcrcus the same information is stored within each uf the follower cache dusters 106 and the leader cache cluster 108 in a number of data shards stored by each of the follower and leader cache clusters 106 and 108, respectively, in the form of a graph including graph nodes and associations or connections between nodes (referred to herein as graph edges!. In particular embodiments, the data shards of each of the follower cache clusters J 06 ami leader cache cluster
CA ? 374065 ΪΟΠ-Ο7-21 tilt: arc buckclizvJ or divided among the cache nodes 112 or I 14 within the respective cache duster, 't hat is, each oi the cache nixies 112 or I 14 within the respective cache cluster stoics a subset of the shanb stored by the cluster (and each sel of shards stored by each ofthe follower and leader cache clusters 106 and 108, respectively, stores the same information, as the lender cache cluster synchronizes the shards stored by each of the cache clusters of a given data center 102, and, in some embodiments, between data centers 102).
In particular embodiments, each graph node is assigned a unique identifier (IO) (hereinafter referred to as node TD) that uniquely identifies the graph node in the graph stored by each ofthe follower and leader cache clusters 106 and 108, respectively, and database 1 10; that
II) is. each node ID is globally unique. In one implementation, each node ID is a M-bit identifier. In one implementation, a shard is allocated a segment of the node ID space. In partictiliir embodiments, each node ID maps (c.g., arithmetically or via conic mathematical function) to a unique corresponding shard ID; UitiL is, each shard ID is also globally unique and refers to the stunc daln object in each set of shards stored by each of tire follower and tender cache clusters
106 mid 108, respectively, in other words, all data objects are stored as graph nodes with unique node IDs and at) the information stored in the graph in the data shards of each of the follower and leader cache clusters 106 and 108. respectively, is stored in the data shards of each of the Ibllowct and tender cectic clusters 106 and 108, respectively, using the same corresponding mil que shard IDs.
As just described, in particular embodiments, the shard ID space (the collection of shard IDs and associated information stored by all the shards of each cache cluster, and replicated in all of tire other follower cadre clusters 106 and leader cache cluster 108) is divided among the follower or leader cache nodes 112 and 114, respectively, within the follower or leader cache clusters 106 and 108, respectively. For example, each follower cache node 112ina given follower cache cluster 106 may store a subset of lltc shards le g., tens, hundreds, or rhtitisunds ot* shards) stored by lire respective follower cache cluster 106 and each shard is assigned u range of node IDs Tor which to store information, including information about the nodes whose respective node IDs map lo lire shard IDs in lltc lunge of shard IDs stored by the particular shard. Similarly, each leader cache node 114 in the leader cache cluster 108 may store a subset of the shards (e.g., tens, hundreds, or thousantls of shards) stored by tire respective
CA 2S74OG5 2017-07-21 leader euche cluster 10K anti each shard is assigned h range of node IDs for which to store information, including in farina lion about the nudes whose respective node IDs map to the shard IDs in the range of shard IDs stored by the particular shard.
However, as described above, a given shard ID corresponds ro Hie same date objects stored by the follower and leader cache clusters !06 and IOS, respectively. As the number of follower cache nodes ί 06 within each follower cache cluster 106 and (lie number of leader cache nodes 114 within the leader cache cluster 108 may vary statically (e.g., the follower cache clusters 106 and ihe leader cache cluster 108 may generally include different numbers of follower cache nodes 112 and leader cache nodes 114. respectively) or dynamically (e.g., cache
It) nodes within a given cache cluster may be shut down for various reasons periodically or ns needed for fixing, updating, or maintenance), the number of shards stoved by each of the follower cache nodes 112 and leader cache nudes 114 may vary statically or dynamically williin each cache cluster as well as between cache clusters. Furthermore. (lie range of shard IDs assigned to eaeh share! may also vary statically or dynamically.
In particular embodiments, each of the follower cache nodes 112 and leader cache nodes 11-1 includes graph management software that manages the storing and serving of information cached within the respective cache node, in particular embodiments, the graph management software running on each of the cache nodes of a given cache cluster may communicate to determine which shards (and corresponding shard IDs) arc stored by each olthe cache nodes within the respective ecchc cluster. Additionally, if the cache node is a follower cache node 112. the graph management software running on the follower cache node 112 receives ictpie.rts (e.g., write or lead requests) from client servers 164, serves the requests by retrieving, up eating, deleting, or storing information within the appropriate shard within the follower caclte node, and manages or facilitates com muni cat ion between the follower cache node
2S 112 and other follower cache nodes I 12 of ihe respective follower cache cluster 106 as well ns communication between the follower cache node 112 and tho lender cache nodes 114 of the lender cache cluster 108. Similarly, if the cache node is a leader cache nude 114, (he graph management software running on the lender cache node 114 manages Ihe communication between the leader cache node 114 and follower cache nodes 112 of the follower cache clusters
106 and the other leader cache nodes 114 of the leader cache cluster 108, as well as
CA 2S74OG5 2017-07-21 eonutionication between the leader cache node 114 and database 110. Tlic graph management software running on each ofthe cache nodes 112 end 114 undetstands that it is storing and serving information in the form of n graph.
in particular embodiments, the graph management software on each follower cache 5 node 112 is also responsible for maintaining a table that it shares with the other cache nodes 112 of tlx; respective follower cache cluster 106. the leader cache nodes 11 4 of the leader cache cluster 108, as well as the client servers 104 dut Ihe respective follower cache cluster 156 serves. This table provides a nitipping of each sliartl ID to ihe particular cache node 112 in a given follower cache cluster 106 that stores the shard SD and information associated with the shard ID.
In this way, the client servers 104 served by a particular follower cache cluster 106 know which of the follower cache nodes 112 within the follower cache cluster 106 maintain the shard ID associated with information the client server 1114 is trying to access, add, or update (e.g., a client server i 04 may send write or read requests to the particular follower cache node 11 2 that stores, or will store, the information associated with a particular shard ID after using the mapping table to determine which of the follower cache nodes ( 12 is assigned, and stores, the short I I Dj. Similarly, in particular embodiments, the graph management software on each leader cache node 114 is also responsible for ruttinUining a table that it shares wilh the other cadre nodes 114 ofthe respective lender cache cluster 10E, as well as (he follower cnehe uodcs 112 of the follower cache dusters 106 (hat the leader cache ci us ter IOS manages. Furthermore, in this way, each follower ccchc node 112 in a given follower cache cluster 106 knows which ofthe other follower cache nodes I 12 in the given follower cache cluster 106 stores which shard IDs stored by the respective follower cache cluster 106. Similarly, in (Iris way each leader cache node I 14 in the lender cache cluster IOS knows which ofthe other leuder cache nodes 114 in the leader cache cluster 108 stores which shnrd IDs stored by the leader cache cluster 108, furthermore, each follower cache node 112 in a given follower cache cLoster 106 knows which of the lender cache nodes 114 in the leader cache cluster 103 stores which shard IDs. Similarly, each leader cache node I 14 iu ihe leader cache cluster I UK knows which of (he follower cache nodes I tl in each of the follower cache clusters 106 stores which shard IDs.
hi particular cm bud i men Is, informal ion regarding each node in the graph, and in particular example embodiments a social graph, is stored in a respective shard of each of (he
CA 2774065 2017-07-21
Follower cache clusters 106 and leader cache clttsicr I OS based on its shard 10. Each node in the graph, as discussed above, has a node ID. Along with the shard ID, the respective cache node 112 ot 114 may store a node type parameter identifying a type of the node, as well as one or more name-value pairs (such as content (e.g., text, media, or URLs to media or other resources})
S and metadata (e.g., a timestamp when the node was created or modified). In particular embodiments, each edge in the graph, and in particular example embodiments a social graph, is stored with each node Lite edge is connected lo. For example, most edges are bi-directional; that is. most edges each connect two nodes in the graph. In particular embodiments, each edge is stored ill the same shard with each node the edge connects. For example, an edge connecting mxle ID 1 to node ID2 may be stored with the shard ID corresponding to node ID I (e.g., shard ID I) and with the sliartl ID corresponding to node 102 (e.g., shard ID2), which may be in dilfcrcnt shards or even di tie rent caclie nodes of a given cache cluster. For example, the edge may be stosed with shard IDl in Ute form of {node ID), edge type, node 1D2) where the edge type indicates lire type of edge. The edge may aiso include metadata (e.g., a timestamp indicating when rile edge was created or modified). The edge may also be cached with shard ID2 in the form of (node IDl, edge type, node ID2). For example, when a user of social networking system 100 establishes a contact relationship with another user or a fan relationship with a enneept or user, the edge relationship of type ‘'friend” or fan’<sup>-</sup> may be stored in two shards, a lityt sinned corresponding lo tire shard to which lire user's identifier is mapped nod a second sherd to which lire object identifier of the other user or concept is mapped.
Networking system 100, and particularly the graph management software running on the follower cache nodes I 12 of follower cache clusters )06 and the leader cache nodes I 14 of the leader cache cluster 108. support a number of queries received from client servais i(M as wefl us to or from other follower or leader cache nodes 112 and ! 14, respectively. For example, the query object add i IDl. node typel. metadata liiot always specified), puylund (not always specified)! causes the receiving cache nude to store a new nude wilh die node tDί specified in the query of the specified node type! in the shard the node IDl corresponds to. The receiving cache notle also stoics with lite node ID J the me lad at a (e.g., a timestamp) and payload (eg, name-value pairs and/or content such us text, media, resources, or references lo resources), if specified As unollter example, the query objccl updatcflDl. node typet (not always specified),
CA 2S7406S SO17-O7-21 ineuidatn [not always specified), payload (not always specified)) causes the receiving cache node lo update Lite node identified by node IDl specified in tlie query (e.g.. change tlw node type to the node type I specified in the query, update the metadata with the met a data specified in the query, or update die content stored with the payload specified in the query) in the corresponding shard. As anotlier example, tlx; query object delete (node TD1 ) causes the receiving cache node tu delete the node identified by node IDl specified in the query. As another example, the query object gct(node IDl J causes the receiving cache node to retrieve the content stored with the node identified by node IDl specified in ihe query.
Now referring to edge queries ftis opposed lo the node queries just described), the
111 query assoc ;uld[int, edge type I, (02. metadata (not always specified)) causes the receiving cache node (which stores node IDl) lo create an edge between the ntxlc identified by node IDl am! (he node identified by node ID2 of edge type edge type 1 and to store (he edge with the node idcniified by node IDl along with (he metadata (e.g., a timestamp indicating when (he edge wns requested) if specified. As another example, the query assoc _updatc(nodc IDl. edge type), node ID2. metadata (not always specified)! causes the receiving cache node (which stores node ID I ) to update (he edge between the node identified by node ID I and the node identified by node iD2. As another example, the query assoc dcletcjnode IDl, edge typcl (not always specified), node ID2} causes ihe receiving cache node (which stores node IDl) to delete the edge between ihe node idcniified by node IDl and the node identified by node 1D2. As another example, the quciy assoc getinodc IDl, edge typel, sottkey (not always specified), start (not always s]5Ceificd), limit (not always specified)} causes die receiving cache node (which stores node 101 ) to return ihe node IDs of the nodes connected to the node identified by node ID I by edges οΓ edge type I, Additionally, if specified, the sottkey specifics a filter. For example, if die sonkey specifies a timestamp, the receiving cache node (which stores node IDl ) returns the node IDs of the nodes connected !o die node identified by node IDl by edges of edge rypcl which were created between (he finie value specified by the slur! parameter and the lime value s]M!cified by the limit parameter. As another example, the query assoc exists {node IDL edge type), list of other node IDs, soit key (not always specified), start (not always specified), limit (not always specified)} causes the receiving cache node (which stores nude IDl) lo return the node IDs ofthe nodes specified in the list of other node IDs connected to tiie node identified by shard IDl by
CA 2Ï74OE5 2017-07-21 edges of edge type I. In addition, the queries described above may be sent in the described form and used to update Ihe leader cache nodes 114.
lti one iinplcmentotion. the caching iuycr implemented by die follower and lcntlcr cnehc clusters IOS and ! Oft cache maintain association data in one or more indexes in a manner that supports high query rates for one or more query types. In some implementations. the invention facilitates efficient intersection. membership and filtering queries directed to associations between nodes in the graph. For example, in One implementation, the caching loyer caches information in a manner optimized to handle point lookup, range and count queries for a variety Of associations between nodes. For example, in constructing a page, a client server 104
If) may issue a query for all friends of a given user. The client server 104 may issue an assoc_gcr query identifying the user and the “friend edge type. To facilitate handling of die query. a cache node in the caching layer may store associations of a given type (such as friends”, funs, 'members, “likes'’, etc.) between a first node (c.fi„ a node corresponding to a user) and a node corresponding to contacts or friends of a user. In addition, to conjunct u not lier party of the page, a client server 104 may issue a query ofthe last N set of wall posts on the profile, by issuing a assoc gel query identifying the user or user profile, the “wallpost” edge type and a limit value. Similarly, comments to a particular wait post can be retrieved in u similar manner.
In one implementation, the caching layer tmpleinenictl by Ihe follower cache clusters 106 nnd the leader cache clusters maintain a set of in-inemory structures for associations between nixies (idl. id2) in tile graph that facilitate fast Searching and handle high query rates. For example, for each fid S .type! associai ion set (a set of all associations that originate at idl and have a given iypcl. the caching layer mainiains nra in-memory indexes. As dimaix^ed above, these· a-auefaiion sets are tin inti med by cache nixies in each duster lhal based on Ihe shard in which id ί falls. Still further, given the structure discussed below, a given association between two iirrttcx may 'v stored in two association acts eneb direelvd io the respective nodes of the nssoeiaiion. Λ fm,s index is based on a teinpoml attribute fs.g., time stamps; and supports range uiterics. λ second index by irJ2 does nor support riijige queries, bid suppure, belter lime • uii'plexriy id'inserts arid look iqw. In une implémentation, the lirai index is an unlured dyiEinnc aiiay of associntii*;; entries stored in a circular buffer. Each entry in the circular buffer describes ur corresponds to one assuuiulioti end eonteiits the following fields; it) iflags 11 byict (iitdiuilirrg
CA 2974065 2017-07-21
Ιΰ tlie visibility olr.ii sssocialmn); M) Sâl2 (K bytes); e) S lime t4 bytes); il) Mata tS by tes J (Sdata is :i fixed κια- S byte It,.Id (wtien more titan * bytes are nct-iled fur StkiLo. this becomes M pointer te unother rneinoty chunk tc hold the ftitl Sdata value. Sdata is opt lui tel foi a given oxsoc type); and e) Slink i.X byicsj offsets of next and previous enti les in the saine h!2 index bucket (see belowi.
lit one ilnplcmcntatiun. the array is ortlered by the Slime attribute ascending. Tits number οΓ entries in the index is capped (such as 10,000) and con figurable by association type. When the limit is readied the array wmpa around, Because the array is Sti me-sorted, most new entries will be appended at the end without shtfnnR any of the existing elements.
In one impie mentation, the primary index can be stored in a single inc me ache key f() that can be looked up by name (assoc:<idI>;<typc>) through a global mcmcached hash table. Hie array can be flunk'd wilh a header containing Ihe following fields: a) count (4 bytes): tire count of visible associations in the (id),type) association set (stored peisistcntly, not just the cached entries in the index); b) head (4 bytes); the byte offset of army farad (dement that sorts highest) in the circular buffer; c) taif (4 bytes) : the byte offset of array tail (element Timt sotts lowest) in tbe circular buffer; and d) id2 index pointer (K bytes): a pointer to a block containing an id? hash table.
The second (Sid?) index is implemented, in one embodiment, us a hush table and supports tjttiek inserts and lookups for a given (Sid!,Srype.Sid2) association. The hash table itself, in one implementation, may be stored in a separate block allocated with memcaehctl’s memory allocator. The table is an array of offsets into the primaiy index, each identifying die fit st element in the corresponding hash bucket. Elements arc linked into a bucket through tlictr Slink fields. Storing the hash table in a separate block allows implcuienters to resize the table and the primary index independently, thus reducing the amount of memory copied as tlie association set grows. Linking association entries into buckets in-place also improves memory efficiency.
The hash table (and bucket lists) may need to bo rebuilt when entries marked hidden or deleted arc expunged from the index, but this can be done infrequently.
Accordingly, as a new associai ion of the same <typc> is added, o cache node I 12, 114 ads the newly associated object lu the hash table and the circular buffer, removing llie oldest entry front the circular buffer. As discussed above, tlie <sortkey> value can be used to sort matching entries based on the attribute, such as a time stamps. In addition, a <liinit<sup>></sup> value limits
CA 29T4O65 2017-07-21 the number of returned results to rite first N values, where N=<linur>. This configuration allows fut setvmg queries regarding associations between nodes at a very high query rate. For example, a first query lusty stsfc to displny a set of friends in a section of a web page. A cache node can quickly respond to a gct^nwoc (idI, type, sortkey, limit) query by looking up association set corrcsjMJiKling to id! by accessing the primary index and retrieving the first N (where N = limit) ic!2 entries in tiie circular buffer. In additiun, the hush table uf the seeuiHlary index facilitates point look ups. Still fuither, the count value maintained by the caching layer facilitates fust responses to the count of a given association set [id I, type).
Some general examples of storing and serving data will now be described (more 10 specific examples relating to particular example implementations of a social graph will be described later alter tiie particular example implementations of the social graph are describcd). t-'or example, when a client server 104 receives a request for a ivcb page, such as from a user of networking system 100, or from another server, component, application, or process οΓ networking system 100 (e.g., in response lo a user request), the client server 104 moy need to issue one or more queries in order to generate the requested web page. In addition, as o user interacts with netwoi king system 100. tiie client server 104 may receive requests that establish or modify object nodes and/or associations be object nodes. It> some instances, the request received by n client server 104 generally includes the node ID representing the user on whose behalf the request to the client server 104 was made. The request may also, or alternately, include one or inure other node IDs corresponding to objects tiie user may want to view, update, delete, or connect or associate (with an edge).
For example, a request may be a read request for accessing information associated with the object or objects die user wants to view (e.g.. one or more objects for serving a web jingo). Fot example, the rend request may be a request for content stored for a particular nude.
For example, a wall post on a user profile can be represented ns a node with an edge type of “wrsllpost. Comments to the wallpost cun also be represented as nodes in the graph with edge type comment'' associations io the wnllpost. In such an example, in particular embodiments, the client server 104 determines the shard ID corresponding to the node ID of the object (node) that includes the content or other information requested, uses the mapping table to determine which of the follower cache nodes 112 (in the follower cache cluster 10ό (hat serves the client
CA 2974065 2017-07-21
IS
Server 104) stoics the shard ID, and transmits a query including the shard ID io tlie particular one of the follower cache nodes 11 2 storing the infonnation associated with and stored with the shard ID. The partieuhir cache node 112 then retrieves the requested information (if cached within the corresponding shard) and transmits the information to the requesting client server 104, which ί may then serve the information to the requesting user (e g., in Ute form of an HTML or other structured document that is rentletablc by the web browser Or other document-rendering application miming on the user’s computing device, if the requested information is not stored/cochcd within the follower cache node 112. the follower cache node 112 tnay then determine, using the mapping table, which of the leader aichc nodes I 14 stores the shaid storing
1(1 the shard ID and forwards the query to the particular leader cache node 114 that stores the shard ID. [(‘the requested information is cached within the partieuhir leader cache node 114, the leader cache node 114 may then retrieve the requested in forma lion tinil forward it to the follower cache i lotie I 12, which ihen updates the particular shard in the follower eaelx; node E12 to store the requested information with tile shard ID and proceeds to serve the query as just described to the client server I (14, which may then serve the information to the requesting user. If the requested information is not caclred within the leader cache node 114, the leader cache node 114 may then translate the query into the language of database 110, and transmit tire new query to tlatabase 110, which then retrieves the requested information and transmits (he requested information lo (he particular leader cache node 114, The lender cache node 114 may then translate lhe retrieved
2(1 io formation hack into the graphical language understood by the graph management software, update the particular shard in the leader cache node 114 to store the requested information with the shard ID, and transmit the retrieved infonnation to the particular follower cache node I 12, which then updates the particular shard in the follower cache node 112 to store the requested information with die shard ID and proceeds to serve the query as just described to the client server 104, which may then serve tire information to the requesting user.
As another example, the user request may be a write request to update existing information or store additional information Tor a node or to create or modify an edge between mo nixies. In lhe former ease, if the iiiftimiiilion to be stored Is for a non-existing node, the client server 104 receiving the user request transmits a request Tor a node ID for a new node lu
3(1 the respective follower cache cluster 106 serving the client server 104. in some cases or
CA 2974065 2017-07-21 embodiments. the client server 104 may specify a prirticulzr shard within which the new node is to be stored (e.g., to eo-kxiile tire new node with another node). In such a case, the client server IÎH requests a new node ID from die pailicular follower cacltc node 112 storing the specified shard- Alternately, the client server 104 may pass a node ID of an existing node with the request for <r new node FD to the follower enelic node I 12 storing the shard that stores the passed node ID to cause the follower cache node 112 to respond to the client server 1 CM with a nude ID fur the new node dial is in the range of node (Ds stored in the shard. In other cases or embodiments, die client server 104 ntay select (e.g., randomly or based on some function) a par tien 1st follower cache node I 12 or a particular shard to send the new node (D request to. Whatever the ease, the particular cache node 112, or more particularly the graph management software running on the follower cache node 112, then transmits the new node ID to the client server 104, The client server 104 may then formulate a «Tile request that includes the new node ID to the corresponding follower cache node 112. The write request may also specify a node type οΓ the new node ami include a payload [e.g., content to be stored with the new node) and/or metadata (e.g-. the node ID of the user making the request, a timestamp indicating when the request was received by the client server 104, among other daw) lo be stored with the node ID. For example, the write request sent to the follower cache node 112 may be uf the form object add {node ID, oode type, payload, metadata}. Similarly, to update a node, the client server 104 may send o write request of the form object modify Inode ID. node type, payload, metadata} to the Follower cache node 112 storing the shtird within which the node ID is stored. Similarly, tn delete a node, the client server 104 may send a request of the form obj ccl_dc let e (node ID} to the follower cache node 1 12 storing the shard within which the shard ID is stored.
In particular embodiments, the follower cache node then transmits the request to the leader cache node 114 storing the shard that stores the eoircspondtng node ID so that the leader cache node 114 may then update the shard. The leader cache node 114 then translates the request into the language of database 110 and transmits the translated request to the database 1 1 0 so that the database may then be updated.
Figure 4 illustrates an example method for processing a request to add an association (assoc add) between two nodes. As Figure 4 illustrates, when a follower cache node
311 112 receives an assoc add request (e.g., axsoc_add(idl, type, id2, metadata), it accesses an index
CA 2974065 2017-07-21 to identify the associalioii set object corresponding to id I and type (402). Follower cache nodes I 12 «dite id2 to both the hush table and die circular buffer of the association set object and increments the count value of the association set object (404). The association set object now maintains the new association of the given type between node id! and node id2. To facilitate searching of the association relative to it12, follower cache node 112 identifies the shard Td corresponding to the nude identifier :d2 and forwards llw assoc_add request to the follower cache node I 12 in the cluster that handles [he identified shard (406). if the instant follower cache node 112 handles the shard, it processes the assoc add request. In one implementation. the forwarding follower cache node 1 12 may (mnsmir a modified assoe_atkl request that signals that
It) this is an update required to establish a bi-directional association in the cache layer. Tire follower enelx.· node 112 tiiso forwards the assoc_atkl request to the leader cache node 114 corresponding lo the shard in which idl falls (408). The leader cache node 114 may execute a similar process to establish a bi-directional association in the leader cache cluster. The lender cache node 114 also causes tire new association to be persisted in database 110. In this in su iter, an association between node idl and node id2 is now searchable in an index with reference to idi nnd type, and separately, itl2 and type.
In particular embodiments. the graph can maintain a variety of different node types, such as users, pages, evenis, wall posts, comments, photographs, videos, background information, concepts, interests and any other element that would be useful to represent as a node. Edge types correspond to associations between (Jtc nodes and can include friends, folio wets, subscribers, fans. likes (or other indications of interest), wallpost, comment, links, suggcsiions, recommendations, and other types of associations between nodes. In one implementation, tt portion of the graph can be a social graph including user nodes that each correspond to a respective user of the social network environment. Tlie social graph may also include other nodes such as concept nodes each devoted or directed to a particular concept as well as topic nodes, which may dt may not be ephemeral. each devoted or directed to a particular topic of cut rent into est among users of the social network environment. In particular embodiments, each node has, represents, or is represented by, a corresponding web page (“profile page”) hosted or accessible in the social network environment. By way of example, a user node may liavc a corresponding user profile page in which lhc corresponding user can add
CA 2974065 2017-07-21 content, make declarations. mid otherwise express himself or herself. By way of example, its will be described below, various web pages hosted or accessible in tlie social network environment such as, for example, user profile pages, concept profile pages, or topic profile pages, enable risers to post content, post status updates, post messages, post comments including comments on other posts submitted by the user or other users, declare interests, declare a “like (described below) towards any of tlie aforementioned posts as well us pages and specific content, or to otherwise express themselves or perform vurious actions (hereinafter these and other user actions limy be collectively referred lo as posts’’ or user actions). In some embodiments, posting may include linking to. or otherwise referencing additional content, such tis media
Id content (e.g., photos, videos, music, text, etc.), uniform resource locators (URLs), and other nodes, via their respective profile pages, other user proClc pages, concept profile pages, topic pages, or oilier web pages or web applications. Such posts, dcelaralions, or actions may then be viewable by the authoring user as well as other users. In particular embodiments, tbe social graph furtlicr includes a plurality of edges that each define or represent a connection between a
IJ corresponding pair of codes in die social graph. As discussed above, cnch item of content may Ik* a node in tbe graph linked to other nodes.
As just described, in various example embodiments, one or more described web pages or web applications arc associated with a social network environment or social networking service. As used herein, a “user may be an individual (human user), an entity (e.g.. an enterprise, business, ot third party application), or a group (e.g., of individuals or eutitics) that interacts or eon tin uni cates with or over such a social network environment. As used herein, a 'registered user” refers to a user that has officially registered within the social network environment (Generally, the users and user nodes described herein refer to registered users only, although ibis is not necessarily a retint re mewl in uther embodiments; that is. in other embodiments, the users uttd user nodes described Iterein may refer to users that have not tegÎstcrcd with the social network environment described herein). In particular embodiments, each user has a corresponding ijroftle page stored, hosted, or accessible by the social network environment and viewable by all or a selected subset of oilier users. Generally, a user has administrative rights to all or a portion of his or her own respective profile page as well as,
3Ü potentially, to other pages created by or for (lie particular user including. for example, home
CA 2374OE5 2017-07-21 pages, pages busting web applications, among Ollier possibilities. As used herein. an “autlienlieatcd user refers to a user who hits been authenticated by the social network environment as being the user claimed in a corresponding profile page to witich the user has administrative rights or, alternately, a suitable hosted representative of the claimed user.
Λ connection between two users or concepts may represent a defined relationship between users or concepts of llic social network environment, and can be defined logically in a suitable data structure of the social network environment as an edge between the nettes corresponding to the users, concepts, events, or oilier nodes of the social network environment fin which the association bus been made. As used herein, a “friendship represents an
I0 association. such as o defined social relationship, between a pair of users of the social network environment. A “friend, as used herein, may refer to any user of the social network environment with which another user has funned a connection, friendship, association, or relationship with, causing an edge to be generated between the two users. Uy way of example, two registered users may become friends with one another explicitly such ns. for example, by one of the two users selecting the other foi friendship as π result of transmitting, or causing to bt transmitted, a friendship request to the other user, who tuny then accept or deny the request. Alternately, friendships or other connections may be automatically established. Such a social friendship maybe visible to other users, especially those who themselves are friends with one or both of the registered ttsers. A friend of a registered user may also have increased access privileges to content, especially user-generated or declared content, on die registered user's profile or other page. It should be noted, however, that two users who have a friend connection established between them in the social graph may not necessarily be friends (in tlx? conventional sense) in real life (outside the social networking environment). For example, in some implementations, a user may be a business or other non-human entity, and thus, incapable of being a friend with a bn man being user in the tradiliutr.il sense οΓ the wont.
As used herein, a “fan'’ may refer lo a user that is a supporter or follower of a particular user, weh page, web application, or other web content accessible in the social nei work environment. In particular embodiments, when a user is a fan of a particular web page (“fans” the particular web page), the user may be listed on that page as a bn for other registered users or the public in general to see. Additionally, an avatar or profile picture of the user may be shown
CA 2Ï74OSS 2017-07-21 on Ihc page (or in/on any of the pages described below). As used herein, a tike may refer to something, such ns, by way of example mid not by way of limitation, a post, a comment, an interest, a link, a piece of media (e.g., photo, photo album, video, song, etc.) a concept, an entity, or a page, among other possibilities (in some implementations a user may indicate or declare a tike to or for virtually anything on any page liostcd by or accessible by the social network system or environment), that a user, and particularly a registered or authenticated user, has declared ot otherwise de mo list titled that he ot she likes, is a fan of, supports, enjoys. Or otherwise has a positive view of. In one embodiment, to indicate or declare a “like or to indicate or declare that the user is a “fan” of something may he processed and defined equivalently in the social networking environment and may be used interchangeably; similarly, to declare oneself a fan” of something, such as a concept ot concept profile page, or lo declare lino oneself “likes” the thing, in:.'y be defined equivalently in the social networking environment anil used interchangeably herein. Additionally, as used herein, an interest may refer to u user-declared interest, such ax a user-declared intetest presented in the user's profile page. As used herein, a “want may refer to virtually anything that u user wants. As described above, a concept may refer io virtually anything that a user may declare or otherwise demonstrate an interest in, a like towards, or a relationship with, such as, by way of example, a sport, a spoils team, a genre of music, a musical composer, a hobby, a business (enterprise), an entity, a group, a celebrity, a person who is nol a registered user, or even, an event, in some embodiments, another user (e.g., a non-auihenticated user), etc. By way of example, there may be a concept node and concept profile page for “Jerry Rice, the famed professional football player, cieutcd and administered by one or more of a plurality of users (e.g., other tlian Jerry Rice), while the social graph additionally includes a user node and nscr profile page for Jeny Rice created by anti administered by Jerry Rice, himself (or trusted or aulhoriicd representatives of Jerry Rice).
Figure 5 illustrates a distributed, redundant system. In the i in piemen ta tion shown, die distributed redundant system includes at least first and second data centers 102a. 102b. Each of the data centers 102a, 102b induites one or more follower cache clusters 106 and a leader eache cluster 108a. 108b. In one implementation. leader cache cluster 108a acts as a primary (master) cuetic cluster, while leader cache cluster 108b is it secondary (slave) cache cluster. I:i one implementation, ilala centers 102a, 102b are redundant ui the sense that synchronization
CA 297*065 2D17-O7-21 functions arc employed to achieve replicated copies of the database 110. In one implementation, data center 102a may be physically located at one geographic region (such as tiie West Coast of the United States) to serve traffic from that region, while data center 102b may be physically located at another geographic region (such as the East Coast of the United States). Given that users from cither of these regions may access the same data and associations, efficient synchronization mechanisms are desired.
Figure 6 illustrates an example method of how a leader cache node 114 proccsscx write commands. As discussed above and with reference to Figure 5, a follower cache node 112 may receive a write command to add/ujxlate an object cr association from a client server KM
I if (Figure 5. No. 1). Tire follower cache node I 12 forwards the write command to a coitospoitding leader cache node 114 (Figure 5. No. 2). When the leader cache node 114 receives a write command from a follower cache node (602), it processes the write command to update oik or more entries in Ihe cache maintained by the leader cadre cluster 108a (604) and writes the update io persistent database 110a (606) (Figure 5. No. 3). The lender cache node 114 also acknowledges tire write command (ACK) to the follower caclic node 1 12 and broadcasts the update to other follower cliche clusters 106 of (lie data center 102a (Figure 5, No. 4a) and the secondary leader cache cluster 108b, which forwards the update to its follower cache clusters 106 (Figure 5, No. 4b) (608). As Figure 6 illustrates, the leader cache node 114 also adds the update to s replication log (6i0), The databases 110a, 110b implement a synchronization mechanism, such os MySQL Replication, to synchronize die persistent databases.
Figure 7 .illustrates a message flow according to one implementation of the invention. When a write command is received at a follower cache node I 12 in a ring 106 thru is not directly associated with the primary leader cache cluster lOBa (Figure 7, No. t ). the follower cache nude 112 forwards the write message to the primary leader cache cluster 108a for processing (Figure 7, No. 2). A leader cache node 114 in the primary leader cadre cluster 108 a may then broadcast the update to its follower cache clusters 106 (Figure 7, No. 3) and writes the changes to database f 10a, As Figure 7 shows, the follower caette node I 12 that received the write command may also forward the write command to ib* secondary lender etiehe cluster 108b (Figure 7, No, 5), which broadcasts the tipilales to other follower cache clusters 106 (Figure 7.
No. 5). The foregoing erctiiledure allows for therefore nllows for changes fo the caching layer fo
CA ÏÏ74OSS £017-07—21 be quickly replicated across data ccmers, while the separate replication between databases I 10a, tOb allow Tor data security.
The applications or processes described herein can be implemented us a scries of computer-readable instructions, embodied or encoded on or within ti tangible data storage medium. that when executed arc operable to cause one or more processors to implement the operations described above. While tire foregoing processes and meclianisms can be implemented by a wide variety of physical systems and in a wide variety of network and computing environments, the computing systems described below provide example computing system architectures of (he server and client systems described above, for didactic, rather than limit: ng, purposes.
Figure 2 illustrates an example computing system architecture, which may be used to implement a server 22a, 22b, In one embodiment, hardware system 1000 comprises a processor 1002, a cache memory 1004, and one or more executable modules and driveis. stored on a tangible computer readable medium, directed to the functions described Iverein.
Additionally, hardware system 1000 includes a high performance inpuj/output (I/O) bus 1006 and a standard I/O bus 1008. Λ hose bridge 1010 couples processor 1002 to high performance I/O bus 1006, whereas I/O bus bridge 1012 couples the two buses 1006 mid 1008 to each other. Λ system memory 1014 and one or more network/coimnumcalion interfaces 1016 couple to bus 1 OOti. I Irirdware system 1000 may further include video memory (not shown) and a display device coupled to the video memory. Mass storage 1018, end I/O ports 1020 couple to bus I0OK. Hardware system 1000 may optionally include a keyboard and pointing device, and a display device (not shown) coupled to bus 1008. Collectively, tircsc elements are intended to represent a bioad category of cuinputcr hardware systems, including but not limited to general purpose computer systems based on the x86-compalible processors manufactured by Intel Corporation of
Santa dura, California, and Ihe x36-compatibie processors manufactured by Advanced Micro Devices ( AMD], lire., ofSuimymlc, California, as well us any other suitable processor.
The elements of hardware system I (W0 are described in greater detail below. In particular, network interlace 1016 provides communication between hardware system 1000 and any of a wide range of networks, such as an Ethernet (e.g,, IEEE 802.3) network, a backplane, etc. Mass storage 1018 provides permanent storage for the data anti programming instructions to
CA 2274OG5 2017-07-21 ixrfoim the above-described functions i in piemen ί cd in the servers 22a, 22b, whereas system memory 1014 (e.g.. DRAM! provides temporary storage for the dam and programming instructions when executed by processor 1002. I/O ports 620 are one or more serial and/or parallel communication ports that provide communication between additional peripheral devices, which may be coupled to hardware system 1000.
Hardware system 1000 may include a variety uf system architectures; and various components of hardware system iUOO may bo rearranged. Fur example, enehc 1004 may be onchip wilh processor 1002. Alternatively, cache 1004 mid processor 1002 may Ire packed together its ίΐ •proccssiw’ module?’ with processor 1002 being referred to as tire “processor cone?’
Furthermore, certain embodiments of the present invention may not require nor include ail of the above components. For example, tbe peripheral devices shown coupled to standard I/O bus IOOS may couple to high performance t/O bus 1006. In addition, in some embodiments, only a single bus may exist, with the components of hardware system 1000 being coirpicd to the single bus. Furthermore, hardware system 1000 may include additional components, such as additional processors, storage devices, or memories.
In one implementation, the operations of the embodiments described herein are implemented us a series of executable modules run by hardware system 1000, individually or collectively in a distributed computing environment. In a particular embodiment, a set of software modules and/or drivers implements a network communications protocol stack, browsing and other computing functions, optimization processes, and the like. The foregoing functional modules may be realized by harthvarc, executable modules stored on a computer readable medium, or a combination of both. For example, tire functional modules may comprise a plurality or series of instructions to be executed by a processor in a hardware system, such as processor HKJ2. initially, the series of instructions may be stored on u storage device, such as mass storage 10 IS. However, the senes of instructions can be tntigibly Stored on any suitable storage medium, such as a diskette, CD-ROM, ROM. EEPROM, etc. Furthermore, the scries of instructions need nut be stored locally, and could be received from a remote storage device, Such as a server on a network, via netwurk/eommunicalions interface 1016. The instructions aie copied from the storage device, such as mass storage 1018, into memory 1014 and then accessed
M and executed by processor 1002.
CA 2774065 ÎO17-O7-21
An operating system moirages ami controls the operation of iuudwaiv system 1000. including the input and output of dam to and from Software applications (not shown). Tire opera ttrg system provides an interface between the software appiications being executed on the system and the hardware components of the system. Any suitublc operating system may be used, such as the LINUX Operating System, the Apple Macintosh Operating System, available from Apple Computer Inc. of Cupertino, Calif,, UNIX operating systems, Microsoft (r) Wtndows(r) operating systems, BSD operating systems, and tlic like. Of course, other implementations aie possible. For example, the nickname generating functions described herein may be implemented ill firmware or on an application specific integrated circuit.
Furthermore, the above-described elements and operations can be comprised of instructions that are stored on storage media. Tlic instructions ean be retrieved and executed by a processing system. Some examples of instructions are software, program code, and firmware. Some examples of storage media urc memory devices, tape, disks, integrated circuits, and servere. The distinctions are operational when executed by the processing system to direct the processing sysiem to operate in accord wilh tlic invention. The term processing system” icfeis to a single processing device or a group of inter-operational processing devices. Some examples of processing devices are integrated circuits and logic circuitry. Those skilled in the art are familiar with instructions, computers, anti storage media.
The present disclosure encompasses ali changes, substitutions, variations, alterations, and modifications ro the example embodiments herein that a person having ordinary skill in the art would comprehend. Similarly, where appropriate, the appended claims encompass all changes, substitutions, variations. alterations, and modifications to the example embodiments herein that u person having ordinary skill in tlic art would comprehend. By way of example, while embodiments of the present invention have been described ns operating in connection with a social networking website, the present invention can be used iu connection with any communications facility that supiwrts web applications and models dota as a graph of associations. Furthermore, in some embodiments the tenu web service and ‘web-site may be used interchangeably and additionally may refer to a custom or generalized API on a device, such ns ;t mobile device (e.g.. cellular phone, smart phone, personal GPS, personal digiial assistance. jx-rsotral gaining device, etc.), that metes API calls directly to a server.
CA 2Ï74OES 2017-07-21
Contents11
6 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6
72 members in 10 offices
Priority claims14
| Document | Office | Kind | Date |
|---|---|---|---|
| 201061428799 | United States of America | P | |
| 201061428799 | United States of America | P | |
| 61428799 | United States of America | – | |
| 13227381 | United States of America | – | |
| 201113227381 | United States of America | A | |
| 201113227381 | United States of America | A | |
| 2964006 | Canada | A | |
| 2964006 | Canada | A | |
| 13227381 | – | – | – |
| 2964006 | – | – | – |
| 61428799 | – | – | – |
| CA20112964006 | – | – | – |
| US201061428799P | – | – | – |
| US201113227381 | – | – | – |
Members72
| Document | Office | Kind | |
|---|---|---|---|
| CA2823187A1 | Canada | A1 | |
| CA2901113A1 | Canada | A1 | |
| CA2911784A1 | Canada | A1 | |
| CA2964006A1 | Canada | A1 | |
| CA2974065A1 | Canada | A1 | |
| US2012173541A1 | United States of America | A1 | |
| US2012173820A1 | United States of America | A1 | |
| US2012173845A1 | United States of America | A1 | |
| WO2012091846A2 | World Intellectual Property Organization (WIPO) | A2 | |
| WO2012091846A3 | World Intellectual Property Organization (WIPO) | A3 | |
| US8438364B2 | United States of America | B2 | |
| AU2011353036A1 | Australia | A1 | |
| CN103380421A | China | A | |
| EP2659386A2 | European Patent Office (EPO) | A2 | |
| MX2013007686A | Mexico | A | |
| US8612688B2 | United States of America | B2 | |
| KR20130143706A | Republic of Korea | A | |
| JP2014501416A | Japan | A | |
| US2014074876A1 | United States of America | A1 | |
| US8832111B2 | United States of America | B2 | |
| US2014330840A1 | United States of America | A1 | |
| US8954675B2 | United States of America | B2 | |
| US2015106359A1 | United States of America | A1 | |
| JP5745649B2 | Japan | B2 | |
| AU2011353036B2 | Australia | B2 | |
| JP2015167034A | Japan | A | |
| AU2015227480A1 | Australia | A1 | |
| CA2823187C | Canada | C | |
| AU2015227480B2 | Australia | B2 | |
| US9208207B2 | United States of America | B2 | |
| CA2901113C | Canada | C | |
| AU2016200251A1 | Australia | A1 | |
| KR101592479B1 | Republic of Korea | B1 | |
| KR101592479B1 | Republic of Korea | B1 | |
| KR20160014111A | Republic of Korea | A | |
| MX337805B | Mexico | B | |
| US2016085881A1 | United States of America | A1 | |
| AU2016200251B2 | Australia | B2 | |
| AU2016203589A1 | Australia | A1 | |
| KR20160083142A | Republic of Korea | A | |
| EP2659386A4 | European Patent Office (EPO) | A4 | |
| KR101640185B1 | Republic of Korea | B1 | |
| CN103380421B | China | B | |
| JP6028065B2 | Japan | B2 | |
| US9514245B2 | United States of America | B2 | |
| JP6062101B1 | Japan | B1 | |
| CN106372136A | China | A | |
| US2017075892A1 | United States of America | A1 | |
| JP2017068852A | Japan | A | |
| JP2017073162A | Japan | A | |
| AU2016203589B2 | Australia | B2 | |
| CA2911784C | Canada | C | |
| AU2017203364A1 | Australia | A1 | |
| KR20170073739A | Republic of Korea | A | |
| MX349037B | Mexico | B | |
| KR101753766B1 | Republic of Korea | B1 | |
| CA2964006C | Canada | C | |
| US9767152B2 | United States of America | B2 | |
| AU2017203364B2 | Australia | B2 | |
| EP2659386B1 | European Patent Office (EPO) | B1 | |
| US9886484B2 | United States of America | B2 | |
| EP3296896A1 | European Patent Office (EPO) | A1 | |
| KR101826115B1 | Republic of Korea | B1 | |
| US2018157660A1 | United States of America | A1 | |
| JP6346255B2 | Japan | B2 | |
| JP2018133100A | Japan | A | |
| CN106372136B | China | B | |
| CA2974065CThis record | Canada | C | |
| US10268725B2 | United States of America | B2 | |
| EP3296896B1 | European Patent Office (EPO) | B1 | |
| JP6584575B2 | Japan | B2 | |
| BR112013016900A2 | Brazil | A2 |
2 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| LapsedLapsedMKLA | MKLA | |
| Examination requestEEER | EEER |
Numbers
- Publication
- 2974065
- Publication, DOCDB
- 2974065
- Publication, EPODOC
- CA2974065
- Application
- 2974065
- Application, DOCDB
- 2974065
- Application, EPODOC
- CA20112974065
Titles2
- English
- DISTRIBUTED CACHE FOR GRAPH DATA
- French
- CACHE REPARTI POUR DONNEES GRAPHIQUES
Classification
- CPC, 16
- G06F16/27
- G06F16/24552
- G06F16/172
- H04L67/568
- G06F16/245
- G06F16/248
- G06F16/284
- G06F16/955
- G06F16/2255
- G06F16/2379
- G06F16/9024
- G06F16/9574
- G06F16/24539
- G06F12/0844
- G06F2212/463
- G06F2212/60
- IPC, 3
- G06F12 0844
- G06F12 0866
- G06F17 30