Use something others have used and configured with documentation and all.
Take a look at distributed files systems. Those usually come with their own share of failure scenarios:
Some I know of are GlusterFS, Ceph (both made by Redhat I think) and Luster.
But as I said now you'll be in the business of understanding and configuring these. Especially how failures can happen.
See first if you can somehow improve the bandwidth between the sites, improve VPN settings, setting QoS flags. Separate networks.
Then take a look of FS-Cache/CacheFS.
http://en.wikipedia.org/wiki/CacheFS
It is a way to use local cache to improve _some_ access patterns. Specifically if you look here:
http://www.linux-mag.com/id/7378/
---
The goal of FS-Cache and CacheFS is to reduce network traffic because some of the data requests will be satisfied by local storage (CacheFS) reducing the amount of network traffic. The load on the server should also be reduced since it will not have to satisfy all data requests. Consequently, this reduction may make up for the increased file lookup time and file read time due to the cache.
---
A key for convenience, which DropBox delivers, is to fetch the data before you need it. 200GB is not that much in the grand scheme of things today - if you only have 1GB/day of changing data, it could be viable.
I've also played with SparkleShare, which has been working well in my very limited testing. The downside is that it is written in .net which has a large footprint (pre-paid if you are in the Microsoft world). With SparkleShare, I'm using my own git repositories (gitolite) which works well but requires some hand editing of the configuration file to make it work.
[1] http://dl.acm.org/ft_gateway.cfm?id=2522721&ftid=1403940&dwn...