Distributed Data Aggregation at Scale for Large Community of Users
Belinda Liu, Thenna Ponnusamy, Adithya Ramakrishnan, Ziyu Lang, Arjun Bhaigond, Amit Desai, Yen Nguyen · 2018
The eCommerce world is facing increasingly huge data volumes and bigger user community. This paper presents an architecture to enable highly performant and highly scalable queries for large community of external customers. The architecture explores the unique pattern of external customer activities: in a big data store hosting a big community of large number of users, in the range of tens or hundreds of millions, each user's data is a fraction of the whole but the community as a whole demands extremely high volume of concurrent analytical queries with sub-second response. In the system, a key-value store is utilized to maximize read concurrency, a custom compression algorithm is developed to minimize data transfer, and a custom query engine is developed to provide aggregation on the fly. Scalability and other potential applications are discussed in the end.