Showing posts with label sharding. Show all posts
Showing posts with label sharding. Show all posts

Mar 2, 2016

Providing MongoDB User Granular Access to User Cluster

Unlike a single instance MongoDB setup or even a ReplicaSet one, when it gets to a Sharded installation, things may get thougher.

For example, if you gave a user a reading permissions to use MongoChef (a most recommended MongoDB client), when it comes to a clustered intallation, in order to avoid the "not authorized to run inprog" error when running db.currentOp(), you should provide the user with some more permissions (in this case the inprog permissions).

Actually it is pretty simple, but it is also a good example for a secured environment management:

Providing inprog Permissions

1. Get to the admin database
use admin; 

2. Authorize as a permitted user
db.auth("admin","admin_password");

3. Create a new role that will have permissions to manage the processes
db.createRole(
{ 
role: "manageOpRole", 
privileges: [ 
{ 
resource: { cluster: true }, 
actions: [ "killop", "inprog" ] 
}, 
{ 
resource: { db: "", collection: "" }, 
actions: [ "killCursors" ] 
} 
], 
roles: [] 
} 
);

4. Provide the permissions to the user:
db.grantRolesToUser(
"reading",
[
      { role: "manageOpRole", db: "admin" }
    ]
);

5. Authenticate as the reading user
db.auth("reading","reading_password");

6. Verify things actually work! (or doing the definition of done);
db.currentOp()

Bottom Line
Simple, tested and secured like we always love our environments!

Keep Performing,
Moshe Kaplan

Feb 21, 2014

When Should I Use MongoDB rather than MySQL (or other RDBMS): The Billing Example

NoSQL is a hot buzz in the air for a pretty long time (well, it not only a buzz anymore).
However, when should we really use it?

Best Practices for MongoDB
NoSQL products (and among them MongoDB) should be used to meet challenges. If you have one of the following challenges, you should consider MongoDB:

You Expect a High Write Load
MongoDB by default prefers high insert rate over transaction safety. If you need to load tons of data lines with a low business value for each one, MongoDB should fit. Don't do that with $1M transactions recording or at least in these cases do it with an extra safety.

You need High Availability in an Unreliable Environment (Cloud and Real Life)
Setting replicaSet (set of servers that act as Master-Slaves) is easy and fast. Moreover, recovery from a node (or a data center) failure is instant, safe and automatic

You need to Grow Big (and Shard Your Data)
Databases scaling is hard (a single MySQL table performance will degrade when crossing the 5-10GB per table). If you need to partition and shard your database, MongoDB has a built in easy solution for that.

Your Data is Location Based
MongoDB has built in spacial functions, so finding relevant data from specific locations is fast and accurate.

Your Data Set is Going to be Big (starting from 1GB) and Schema is Not Stable
Adding new columns to RDBMS can lock the entire database in some database, or create a major load and performance degradation in other. Usually it happens when table size is larger than 1GB (and can be major pain for a system like BillRun that is described bellow and has several TB in a single table). As MongoDB is schema-less, adding a new field, does not effect old rows (or documents) and will be instant. Other plus is that you do not need a DBA to modify your schema when application changes.


You Don't have a DBA
If you don't have a DBA, and you don't want to normalize your data and do joins, you should consider MongoDB. MongoDB is great for class persistence, as classes can be serialized to JSON and stored AS IS in MongoDB. Note: If you are expecting to go big, please notice that your will need to follow some best practices to avoid pitfalls


Real World Case Study: Billing
In the last ILMUG, Ofer Cohen presented BillRun, a next generation Open Source billing solution that utilizes MongoDB as its data store. This billing system runs in production in the fastest growing cellular operator in Israel, where it processes over 500M CDRs (call data records) each month. In his presentation Ofer presented how this system utilizes MongoDB advantages:
  1. Schema-less design enables rapid introduction of new CDR types to the system. It let BillRun keep the data store generic.
  2. Scale BillRun production site already manages several TB in a single table, w/o being limited by adding new fields or being limited by growth
  3. Rapid replicaSet enables meeting regulation with easy to setup multi data center DRP and HA solution.
  4. Sharding enables linear and scale out growth w/o running out of budget.
  5. With over 2,000/s CDR inserts, MongoDB architecture is great for a system that must support high insert load. Yet you can guarantee transactions with findAndModify (which is slower) and two-phase commit (application wise).
  6. Developer oriented queries, enable developers write a elegant queries.
  7. Location based is being utilized to analyze users usage and determining where to invest in cellular infrastructure.
Bottom Line
MongoDB is great tool, that should be used in the right scenarios to gain unfair advantage in your market. BillRun is a fine example for that.

Keep Performing,
Moshe Kaplan

Aug 31, 2011

Sharding is Now COTS

I wrote here a lot in the past regarding SQL Server and MySQL sharding. I wanted to update you that ScaleBase, which is led by the industry veterans Liran Zelkha and Doron Levari announced that ScaleBase 1.0 is now available. 
Scalebase provides a fine solution to MySQL sharding by implementing an SQL load balancer in front of your servers. 
This product is available as a service and as a downloadable product. ScaleBase 1.0 supports:
  1. Read/Write splitting
  2. Transparent Sharding
  3. High Availability
If you need a proper solution for a large size MySQL database, yet you don't have the time or mind to implement sharding by yourselves, you should take a look at their product.


Keep Performing,
Moshe Kaplan Follow MosheKaplan on Twitter

Mar 28, 2011

MySQL Partitioning. One step before Sharding.

Many of you may consider Sharding to meet your large scale database demands. However, in some cases it seems too early since you may not have the needed bandwidth in the current phase or you just consider to adapt a Sharding COTS solution like Scalebase and you need a solution for mid term.


I don't have time so what should I do?
The answer is simple: you may choose MySQL built-in mechanism named MySQL Partitioning. 
This mechanism partitions your table using one of 4 different rule types. It main solution is keeping all this process hidden from the applicative DBA and the software engineer. Try one of this methods:

  1. Range Partitioning: partition by column values. Most fit for known ranges like archive by date range.
  2. List Partitioning: similar to the above based on discrete values. Most fit for archive by years.
  3. Hash Partitioning: partition by hash function defined by user. Most fit when data ranges are that are unknown, but you know the spread of them. Should be used when you partition a table according to a foreign key. Please notice that in any case, the hashed key must be included in table primary key.
  4. Key Partitioning: similar to Hash, but this time the MySQL Server takes care of the function. Most fit that the data distribution is uniform. Partition key should include the primary key.
How do I declare that?
Using DDL. just add few more lines to your CREATE TABLE statement and you will receive the partitioning. For example to split by year use the following statement:

CREATE TABLE grades (
...    
year INT NOT NULL,
...
)
PARTITION BY LIST(year) (
    PARTITION p2009 VALUES IN (2009),
    PARTITION p2010 VALUES IN (2010),
    PARTITION p2011 VALUES IN (2011),
    PARTITION p2012 VALUES IN (2012)
);

What do I get?
Mayflower has presented very nice numbers with 200M records table partitioned to 400 parts. They reached 4000 INSERT/UPDATE and 6000 SELECT statements per second.

Bottom Line
If don't have a lot of time and you need an instant solution, go for Partitioning

Keep Performing,

Dec 31, 2010

Prepare for Database2011

Hi,


Raphael Fogel invited me to give a lecture at Database2011, the central databases conference in Israel. The event will take place on Jan 13 at the Avenue Conference Center.
So reserve the date and come prepare to hear How Sharding turned MySQL into the Internet de-facto database standard? What is Sharding? Why the biggest internet players chose MySQL? and What are the latest solution in this field.


Keep Performing and Happy Holidays,
Moshe KaplanFollow MosheKaplan on Twitter

Feb 21, 2010

Lecture: Extract The Traffic from the DB

A few days ago I had a presentation in the AlphaGeeks meetup in Tel Aviv, presenting NoSQL, Memcached, CouchDB, Sharding and other buzzwords that help you extract the traffic from the database and boost your system performance. If you were not there, you can take a look at the presentation (English) or the recorded video (Hebrew).

The Presentation (English)
The Video (Hebrew)

Keep Performing,
Moshe Kaplan

Jun 22, 2009

Billion Events per Day, Israel 3rd Java Technology Day, June 22, 2009

We presented today at the Israeli 3rd Java Technology Day, the largest SUN Microsystems/MySQL event in Israel. We presented here the essentials parts of building a real life web/enterprise system that needs to handle the performance needs of 1 billion events per day (a case study from the ad networks billing systems). We presented the adoption rate in the internet, Load Balancers (HAProxy, Apache, Radware, F5, Cisco), Web Servers, In Memory Database (IMDB inc. Memcached, Gigaspaces, Teracotta and Oracle Coherence) and finally Sharding (inc. Veritical, Static Horizontal and dynamic). A great example for a performance boosting architecture.

Feel free to take a look at the presentation:


Apr 22, 2009

Very large databases on the cloud: The Presentation

This time I include the presentation that I presented on Monday: "How your very large databases can work in the cloud computing world?". It was a great presentation and the place was loaded with industry experts from various companies such as Amdocs, Panaya, SAP, Superfish and SUN.

The presentation content can be found bellow as well as the slide show:

Cloud computing is famous for its flexibility, dynamic nature and ability to infinite growth. However, infinite growth means very large databases with billions of records in it. This leads us to a paradox: "How can weak servers support very large databases which usually require several CPUs and dedicated hardware?"
The Internet industry proved it can be done. These days many of the Internet giants, processing billions of events every day, are based on cloud computing architecture such and sharding. What is Sharding? What kinds of Sharding can you implement? What are the best practices?

Keep Performing,
Moshe Kaplan. RockeTier. The Performance and Cloud Experts.

Apr 18, 2009

Very large databases on the cloud

This Monday, we'll present in the IGT cloud computing workgroup: "How your very large databases can work in the cloud computing world?". The presentation will be held along with other presentations by Nati Shalom and Haim Yadid, market experts in the field of performance and cloud computing. Therefore, it will be interesting being there.

How your very large databases can work in the cloud computing world?
Moshe Kaplan, RockeTier, a performance expert and scale out architect
Cloud computing is famous for its flexibility, dynamic nature and ability to infinite growth. However, infinite growth means very large databases with billions of records in it. This leads us to a paradox: "How can weak servers support very large databases which usually require several CPUs and dedicated hardware?"
The Internet industry proved it can be done. These days many of the Internet giants, processing billions of events every day, are based on cloud computing architecture such and sharding. What is Sharding? What kinds of Sharding can you implement? What are the best practices?

Date: Apr 20, 2009 14:00-17:00
Location: IGT Offices, Maskit 4, Hertzelia Ind. Zone, Israel.
Confirmation at: info@grid.org.il

Keep Performing,
Moshe Kaplan. RockeTier. The Performance and Cloud Experts.

Apr 14, 2009

MySQL Sharding

A few weeks ago we had a presentation in the Israeli MySQL User Group, where we presented "How Sharding turned MySQL into the Internet de-facto database standard?"
This presentation dealt with the common belief in the enterprise software world that MySQL cannot scale to large databases sizes. The Internet industry proved it can be done. These days many of the Internet giants, processing billions of events every day, are based on MySQL. Most of these giants were able to turn MySQL into a mighty database machines by implementing Sharding.
In the attached presentation from SlideShare we answer the following questions: What is Sharding? What kinds of Sharding can you implement? What are the best practices?



Keep Performing,
Moshe Kaplan. RockeTier. The Performance Experts.

Mar 18, 2009

The eBay way

Hi,

I received from Shachar Zehavi, our new director of R&D, a link to Chris Kasten's presentation regarding eBay

In this presentation Chris presents a simple but very clever solution to support 4 billion events per day on 25 commodity servers. Why instead of using complex in memory databases solutions? not using MySQL as the ultimate grid solution?

What are the key components in eBay solution:
1. In memory database: implementation using MySQL in memory engine
2. HA: implemented using MySQL replication between two different MySQL in memory databases
3. Persistence: implementation using batch process once in 5min. In this case the store is done using InnoDB one.
4. Scalability: can b achieved using horizontal sharding

The number described in this presentation (4 billion requests per day using 25 machines) remind me another presenation of Paul Strong, distinguished research scientist from eBay. This presensation at the IGT2008 - The World Summit of Cloud Computing described eBay numbers and challenges including: 150 Billion request per day which is about 200K requests per second, which is a remarkable number.

Keep Performing,
Moshe Kaplan. RockeTier. The Performance Experts.

Feb 19, 2009

Lecture: MySQL Sharding

Hi,

In the next Israel MySQL User Group, RockeTier will present:

"How Sharding turned MySQL into the Internet de-facto Database Standard?"
Abstract
A common belief in the enterprise software world is that MySQL cannot scale to large databases sizes. The Internet industry proved it can be done. These days many of the Internet giants, processing billions of events every day, are based on MySQL. Most of these giants were able to turn MySQL into a mighty database machines by implementing Sharding.
What is Sharding? What kinds of Sharding can you implement? What are the best practices? All these issues will be address in this lecture by Moshe Kaplan from RockeTier, a performance expert and scale out architect

When: Wed, March 4th
Where: InterBit, 6 Ha`chilazon St., Ramat Gan, Israel, 03-7529922

Jan 20, 2009

Does MySQL 5.0 work with multi-core processors? yes???


Hi,

I think this week can truly be named: "The MySQL week". So many issues regarding this product. You know, sometimes it is the cost of getting things for free...

Well, first of all, MySQL is a great product. Many start ups started with this product and many giants are using it for their billions events per day systems.

However, MySQL has several limitations. One of them is that it not really supports multi core processors. Yes, I know, MySQL definitly tell that they do support multi threading. However, our analysis found out that this is not the case. As you can see in the attached graph, the MySQL machine really works hard. However, since it's a quad core machine, it is reaching only 25% CPU utilization

To be more accorate, too many people around the globe like Jeremy Kusnetz, peter with who Sun reach performance on 256 cores and starnixhacks got to the same conclusion: MySQL do use multi threading to it periphrial components, but when it get to do real work, it not really supports multi core processors.

This can be really defined as a major issue, since modern CPUs are based on slower clocks and more cores... So what can we do?

Well the answer is Sharding...
In few words, sharing is the internet companies way to install many weak type databases, each dealing with a vertical or horizental part of the database. This way you can install on a dual qual core CPUs machine, 8 MySQLs instances, each dealing with a single partition of your application (for example the first stores customer with name starts with 'A', while the second stores these that start with 'B', 'C' and 'D' and so on). A broader review of this method will be provided in the next weeks.

Anyhow, I recommend MySQL Optimization guide, as a must to any MySQL would be tuning professional

Moshe.
RockeTier, The Performance Experts

Nov 24, 2008

Online advertisement. How do you handle billion events per day?

Hi,

It was an interesting day today. I'm taking a part in a two days conference named Affilicon. This is the first affiliates networks and affiliates conference in Israel.

Why do I care?
Well, we have several clients in this field.

Why these clients are interested in our services?
Well, it seems that online advertisement systems is a major source for high end system that must handle thousands of events per seconds, or billions of events per day!

Are there so many?
Ya, these systems are counting every banner, text ad or video impression. Think for a second about Google Adwords. Have you taken a look on your conversation rates? Did you think how many text ads do they serve per day? Did you think how do they count all these impressions?

How do they handle these rates?
Well, these players have very large farms (Google for example has around 1 Millions servers in its data centers). They also implement complex solutions to handle these stress in complex ways, doing it all using commodity servers

So how do RockeTier help them?
Well in various ways: 1) boosting their performance, reducing their number of servers by factor of up to 20; 2) implementing shredding, load balancers and grid solutions in order to split the processing between several servers and reducing the stress on each server and 3) use better algorithms to better their system

Best,
Moshe
RockeTier

ShareThis

Intense Debate Comments

Ratings and Recommendations