merlimat closed pull request #1764: Pulsar Go client library
URL: https://github.com/apache/incubator-pulsar/pull/1764
 
 
   

This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:

As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):

diff --git a/pom.xml b/pom.xml
index fa3d4b76c3..39a31de5ed 100644
--- a/pom.xml
+++ b/pom.xml
@@ -953,9 +953,11 @@ flexible messaging model and an intuitive client 
API.</description>
             <exclude>site/img/**</exclude>
             <exclude>generated-site/**</exclude>
             <exclude>.github/*.md</exclude>
+            <exclude>**/.idea/*</exclude>
           </excludes>
           <mapping>
             <proto>JAVADOC_STYLE</proto>
+            <go>DOUBLESLASH_STYLE</go>
             <conf>SCRIPT_STYLE</conf>
             <ini>SCRIPT_STYLE</ini>
             <yaml>SCRIPT_STYLE</yaml>
diff --git a/pulsar-client-cpp/include/pulsar/c/client_configuration.h 
b/pulsar-client-cpp/include/pulsar/c/client_configuration.h
index b04c21d03a..d1979c480b 100644
--- a/pulsar-client-cpp/include/pulsar/c/client_configuration.h
+++ b/pulsar-client-cpp/include/pulsar/c/client_configuration.h
@@ -27,7 +27,8 @@ extern "C" {
 
 typedef enum { pulsar_DEBUG = 0, pulsar_INFO = 1, pulsar_WARN = 2, 
pulsar_ERROR = 3 } pulsar_logger_level_t;
 
-typedef void (*pulsar_logger)(pulsar_logger_level_t level, const char *file, 
int line, const char *message);
+typedef void (*pulsar_logger)(pulsar_logger_level_t level, const char *file, 
int line, const char *message,
+                              void *ctx);
 
 typedef struct _pulsar_client_configuration pulsar_client_configuration_t;
 typedef struct _pulsar_authentication pulsar_authentication_t;
@@ -105,7 +106,8 @@ void 
pulsar_client_configuration_set_concurrent_lookup_request(pulsar_client_con
  */
 int 
pulsar_client_configuration_get_concurrent_lookup_request(pulsar_client_configuration_t
 *conf);
 
-void pulsar_client_configuration_logger(pulsar_client_configuration_t *conf, 
pulsar_logger logger);
+void pulsar_client_configuration_set_logger(pulsar_client_configuration_t 
*conf, pulsar_logger logger,
+                                            void *ctx);
 
 void pulsar_client_configuration_set_use_tls(pulsar_client_configuration_t 
*conf, int useTls);
 
diff --git a/pulsar-client-cpp/include/pulsar/c/message.h 
b/pulsar-client-cpp/include/pulsar/c/message.h
index 7c2fb9dce5..107fe6cb7a 100644
--- a/pulsar-client-cpp/include/pulsar/c/message.h
+++ b/pulsar-client-cpp/include/pulsar/c/message.h
@@ -26,6 +26,8 @@ extern "C" {
 #include <stddef.h>
 #include <stdint.h>
 
+#include "string_map.h"
+
 #pragma GCC visibility push(default)
 
 typedef struct _pulsar_message pulsar_message_t;
@@ -102,7 +104,7 @@ void pulsar_message_disable_replication(pulsar_message_t 
*message, int flag);
  *
  * @return an unmodifiable view of the properties map
  */
-// const StringMap& getProperties() const;
+pulsar_string_map_t *pulsar_message_get_properties(pulsar_message_t *message);
 
 /**
  * Check whether the message has a specific property attached.
diff --git a/pulsar-client-cpp/include/pulsar/c/string_map.h 
b/pulsar-client-cpp/include/pulsar/c/string_map.h
new file mode 100644
index 0000000000..b3c2188aeb
--- /dev/null
+++ b/pulsar-client-cpp/include/pulsar/c/string_map.h
@@ -0,0 +1,46 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+#pragma once
+
+#ifdef __cplusplus
+extern "C" {
+#endif
+
+#pragma GCC visibility push(default)
+
+typedef struct _pulsar_string_map pulsar_string_map_t;
+
+pulsar_string_map_t *pulsar_string_map_create();
+void pulsar_string_map_free(pulsar_string_map_t *map);
+
+int pulsar_string_map_size(pulsar_string_map_t *map);
+
+void pulsar_string_map_put(pulsar_string_map_t *map, const char *key, const 
char *value);
+
+const char *pulsar_string_map_get(pulsar_string_map_t *map, const char *key);
+
+const char *pulsar_string_map_get_key(pulsar_string_map_t *map, int idx);
+const char *pulsar_string_map_get_value(pulsar_string_map_t *map, int idx);
+
+#pragma GCC visibility pop
+
+#ifdef __cplusplus
+}
+#endif
\ No newline at end of file
diff --git a/pulsar-client-cpp/lib/ClientConfigurationImpl.h 
b/pulsar-client-cpp/lib/ClientConfigurationImpl.h
index b5da7908ed..7160a6592b 100644
--- a/pulsar-client-cpp/lib/ClientConfigurationImpl.h
+++ b/pulsar-client-cpp/lib/ClientConfigurationImpl.h
@@ -45,8 +45,8 @@ struct ClientConfigurationImpl {
           logConfFilePath(),
           useTls(false),
           tlsAllowInsecureConnection(true),
-          statsIntervalInSeconds(600) {  // 10 minutes
-    }
+          statsIntervalInSeconds(600),  // 10 minutes
+          loggerFactory() {}
 };
 }  // namespace pulsar
 
diff --git a/pulsar-client-cpp/lib/LogUtils.cc 
b/pulsar-client-cpp/lib/LogUtils.cc
index 219232702d..e2615a54e6 100644
--- a/pulsar-client-cpp/lib/LogUtils.cc
+++ b/pulsar-client-cpp/lib/LogUtils.cc
@@ -25,7 +25,7 @@
 
 namespace pulsar {
 
-void LogUtils::init(const std::string &logfilePath) {
+void LogUtils::init(const std::string& logfilePath) {
 // If this is called explicitely, we fallback to Log4cxx config, if enabled
 
 #ifdef USE_LOG4CXX
@@ -48,4 +48,11 @@ LoggerFactoryPtr LogUtils::getLoggerFactory() {
     return s_loggerFactory;
 }
 
+std::string LogUtils::getLoggerName(const std::string& path) {
+    // Remove all directories from filename
+    int startIdx = path.find_last_of("/");
+    int endIdx = path.find_last_of(".");
+    return path.substr(startIdx + 1, endIdx - startIdx - 1);
+}
+
 }  // namespace pulsar
\ No newline at end of file
diff --git a/pulsar-client-cpp/lib/LogUtils.h b/pulsar-client-cpp/lib/LogUtils.h
index 12a5b0696c..81de443b46 100644
--- a/pulsar-client-cpp/lib/LogUtils.h
+++ b/pulsar-client-cpp/lib/LogUtils.h
@@ -29,15 +29,16 @@ namespace pulsar {
 
 #define PULSAR_UNLIKELY(expr) __builtin_expect(expr, 0)
 
-#define DECLARE_LOG_OBJECT()                                                   
                    \
-    static pulsar::Logger* logger() {                                          
                    \
-        static boost::thread_specific_ptr<pulsar::Logger> 
threadSpecificLogPtr;                    \
-        pulsar::Logger* ptr = threadSpecificLogPtr.get();                      
                    \
-        if (PULSAR_UNLIKELY(!ptr)) {                                           
                    \
-            
threadSpecificLogPtr.reset(pulsar::LogUtils::getLoggerFactory()->getLogger(__FILE__));
 \
-            ptr = threadSpecificLogPtr.get();                                  
                    \
-        }                                                                      
                    \
-        return ptr;                                                            
                    \
+#define DECLARE_LOG_OBJECT()                                                   
                  \
+    static pulsar::Logger* logger() {                                          
                  \
+        static boost::thread_specific_ptr<pulsar::Logger> 
threadSpecificLogPtr;                  \
+        pulsar::Logger* ptr = threadSpecificLogPtr.get();                      
                  \
+        if (PULSAR_UNLIKELY(!ptr)) {                                           
                  \
+            std::string logger = pulsar::LogUtils::getLoggerName(__FILE__);    
                  \
+            
threadSpecificLogPtr.reset(pulsar::LogUtils::getLoggerFactory()->getLogger(logger));
 \
+            ptr = threadSpecificLogPtr.get();                                  
                  \
+        }                                                                      
                  \
+        return ptr;                                                            
                  \
     }
 
 #define LOG_DEBUG(message)                                                 \
@@ -85,6 +86,8 @@ class LogUtils {
     static void setLoggerFactory(LoggerFactoryPtr loggerFactory);
 
     static LoggerFactoryPtr getLoggerFactory();
+
+    static std::string getLoggerName(const std::string& path);
 };
 
 #pragma GCC visibility pop
diff --git a/pulsar-client-cpp/lib/SimpleLoggerImpl.cc 
b/pulsar-client-cpp/lib/SimpleLoggerImpl.cc
index 309eb33da1..95b2585995 100644
--- a/pulsar-client-cpp/lib/SimpleLoggerImpl.cc
+++ b/pulsar-client-cpp/lib/SimpleLoggerImpl.cc
@@ -78,13 +78,7 @@ class SimpleLogger : public Logger {
     }
 };
 
-Logger *SimpleLoggerFactory::getLogger(const std::string &path) {
-    // Remove all directories from filename
-    int startIdx = path.find_last_of("/");
-    int endIdx = path.find_last_of(".");
-    std::string fileName = path.substr(startIdx + 1, endIdx - startIdx - 1);
-    return new SimpleLogger(fileName);
-}
+Logger *SimpleLoggerFactory::getLogger(const std::string &file) { return new 
SimpleLogger(file); }
 
 LoggerFactoryPtr SimpleLoggerFactory::create() { return LoggerFactoryPtr(new 
SimpleLoggerFactory); }
 }  // namespace pulsar
diff --git a/pulsar-client-cpp/lib/c/cStringMap.cc 
b/pulsar-client-cpp/lib/c/cStringMap.cc
new file mode 100644
index 0000000000..221dce4b57
--- /dev/null
+++ b/pulsar-client-cpp/lib/c/cStringMap.cc
@@ -0,0 +1,60 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+#include <pulsar/c/string_map.h>
+
+#include "c_structs.h"
+
+pulsar_string_map_t *pulsar_string_map_create() { return new 
pulsar_string_map_t; }
+
+void pulsar_string_map_free(pulsar_string_map_t *map) { delete map; }
+
+int pulsar_string_map_size(pulsar_string_map_t *map) { return map->map.size(); 
}
+
+void pulsar_string_map_put(pulsar_string_map_t *map, const char *key, const 
char *value) {
+    map->map[key] = value;
+}
+
+const char *pulsar_string_map_get(pulsar_string_map_t *map, const char *key) {
+    std::map<std::string, std::string>::iterator it = map->map.find(key);
+
+    if (it == map->map.end()) {
+        return NULL;
+    } else {
+        return it->second.c_str();
+    }
+}
+
+const char *pulsar_string_map_get_key(pulsar_string_map_t *map, int idx) {
+    std::map<std::string, std::string>::iterator it = map->map.begin();
+    while (idx-- > 0) {
+        ++it;
+    }
+
+    return it->first.c_str();
+}
+
+const char *pulsar_string_map_get_value(pulsar_string_map_t *map, int idx) {
+    std::map<std::string, std::string>::iterator it = map->map.begin();
+    while (idx-- > 0) {
+        ++it;
+    }
+
+    return it->second.c_str();
+}
\ No newline at end of file
diff --git a/pulsar-client-cpp/lib/c/c_Client.cc 
b/pulsar-client-cpp/lib/c/c_Client.cc
index 1063bb877d..905e41067f 100644
--- a/pulsar-client-cpp/lib/c/c_Client.cc
+++ b/pulsar-client-cpp/lib/c/c_Client.cc
@@ -26,7 +26,7 @@
 pulsar_client_t *pulsar_client_create(const char *serviceUrl,
                                       const pulsar_client_configuration_t 
*clientConfiguration) {
     pulsar_client_t *c_client = new pulsar_client_t;
-    c_client->client.reset(new pulsar::Client(std::string(serviceUrl)));
+    c_client->client.reset(new pulsar::Client(std::string(serviceUrl), 
clientConfiguration->conf));
     return c_client;
 }
 
diff --git a/pulsar-client-cpp/lib/c/c_ClientConfiguration.cc 
b/pulsar-client-cpp/lib/c/c_ClientConfiguration.cc
index 6ceaf6b103..935e908530 100644
--- a/pulsar-client-cpp/lib/c/c_ClientConfiguration.cc
+++ b/pulsar-client-cpp/lib/c/c_ClientConfiguration.cc
@@ -72,28 +72,34 @@ int 
pulsar_client_configuration_get_concurrent_lookup_request(pulsar_client_conf
 class PulsarCLogger : public pulsar::Logger {
     std::string file_;
     pulsar_logger logger_;
+    void *ctx_;
 
    public:
-    PulsarCLogger(const std::string &file, pulsar_logger logger) : 
file_(file), logger_(logger) {}
+    PulsarCLogger(const std::string &file, pulsar_logger logger, void *ctx)
+        : file_(file), logger_(logger), ctx_(ctx) {}
 
     bool isEnabled(Level level) { return level >= pulsar::Logger::INFO; }
 
     void log(Level level, int line, const std::string &message) {
-        logger_((pulsar_logger_level_t)level, file_.c_str(), line, 
message.c_str());
+        logger_((pulsar_logger_level_t)level, file_.c_str(), line, 
message.c_str(), ctx_);
     }
 };
 
 class PulsarCLoggerFactory : public pulsar::LoggerFactory {
     pulsar_logger logger_;
+    void *ctx_;
 
    public:
-    PulsarCLoggerFactory(pulsar_logger logger) : logger_(logger) {}
+    PulsarCLoggerFactory(pulsar_logger logger, void *ctx) : logger_(logger), 
ctx_(ctx) {}
 
-    pulsar::Logger *getLogger(const std::string &fileName) { return new 
PulsarCLogger(fileName, logger_); }
+    pulsar::Logger *getLogger(const std::string &fileName) {
+        return new PulsarCLogger(fileName, logger_, ctx_);
+    }
 };
 
-void pulsar_client_configuration_set_logger(pulsar_client_configuration_t 
*conf, pulsar_logger logger) {
-    conf->conf.setLogger(pulsar::LoggerFactoryPtr(new 
PulsarCLoggerFactory(logger)));
+void pulsar_client_configuration_set_logger(pulsar_client_configuration_t 
*conf, pulsar_logger logger,
+                                            void *ctx) {
+    conf->conf.setLogger(pulsar::LoggerFactoryPtr(new 
PulsarCLoggerFactory(logger, ctx)));
 }
 
 void pulsar_client_configuration_set_use_tls(pulsar_client_configuration_t 
*conf, int useTls) {
diff --git a/pulsar-client-cpp/lib/c/c_Message.cc 
b/pulsar-client-cpp/lib/c/c_Message.cc
index d87560e374..f9288ac863 100644
--- a/pulsar-client-cpp/lib/c/c_Message.cc
+++ b/pulsar-client-cpp/lib/c/c_Message.cc
@@ -94,3 +94,9 @@ uint64_t 
pulsar_message_get_publish_timestamp(pulsar_message_t *message) {
 uint64_t pulsar_message_get_event_timestamp(pulsar_message_t *message) {
     return message->message.getEventTimestamp();
 }
+
+pulsar_string_map_t *pulsar_message_get_properties(pulsar_message_t *message) {
+    pulsar_string_map_t *map = pulsar_string_map_create();
+    map->map = message->message.getProperties();
+    return map;
+}
diff --git a/pulsar-client-cpp/lib/c/c_structs.h 
b/pulsar-client-cpp/lib/c/c_structs.h
index a4ff193ca2..41e9fba668 100644
--- a/pulsar-client-cpp/lib/c/c_structs.h
+++ b/pulsar-client-cpp/lib/c/c_structs.h
@@ -79,4 +79,8 @@ static void handle_result_callback(pulsar::Result result, 
pulsar_result_callback
     if (callback) {
         callback((pulsar_result)result, ctx);
     }
-}
\ No newline at end of file
+}
+
+struct _pulsar_string_map {
+    std::map<std::string, std::string> map;
+};
diff --git a/pulsar-client-go/examples/consumer-listener/consumer-listener.go 
b/pulsar-client-go/examples/consumer-listener/consumer-listener.go
new file mode 100644
index 0000000000..8e7962c4fa
--- /dev/null
+++ b/pulsar-client-go/examples/consumer-listener/consumer-listener.go
@@ -0,0 +1,60 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package main
+
+import (
+       "github.com/apache/incubator-pulsar/pulsar-client-go/pulsar"
+       "fmt"
+       "log"
+)
+
+func main() {
+       client, err := pulsar.NewClient(pulsar.ClientOptions{URL: 
"pulsar://localhost:6650"})
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer client.Close()
+
+       channel := make(chan pulsar.ConsumerMessage)
+
+       consumer, err := client.Subscribe(pulsar.ConsumerOptions{
+               Topic:            "my-topic",
+               SubscriptionName: "my-subscription",
+               Type:             pulsar.Shared,
+               MessageChannel:   channel,
+       })
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer consumer.Close()
+
+       // Receive messages from channel. The channel returns a struct which 
contains message and the consumer from where
+       // the message was received. It's not necessary here since we have 1 
single consumer, but the channel could be
+       // shared across multiple consumers as well
+       for cm := range channel {
+               msg := cm.Message
+               fmt.Printf("Received message  msgId: %s -- content: '%s'\n",
+                       msg.ID(), string(msg.Payload()))
+
+               consumer.Ack(msg)
+       }
+}
diff --git a/pulsar-client-go/examples/consumer/consumer.go 
b/pulsar-client-go/examples/consumer/consumer.go
new file mode 100644
index 0000000000..7a10128547
--- /dev/null
+++ b/pulsar-client-go/examples/consumer/consumer.go
@@ -0,0 +1,59 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package main
+
+import (
+       "github.com/apache/incubator-pulsar/pulsar-client-go/pulsar"
+       "fmt"
+       "log"
+       "context"
+)
+
+func main() {
+       client, err := pulsar.NewClient(pulsar.ClientOptions{URL: 
"pulsar://localhost:6650"})
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer client.Close()
+
+       consumer, err := client.Subscribe(pulsar.ConsumerOptions{
+               Topic:            "my-topic",
+               SubscriptionName: "my-subscription",
+               Type:             pulsar.Shared,
+       })
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer consumer.Close()
+
+       for {
+               msg, err := consumer.Receive(context.Background())
+               if err != nil {
+                       log.Fatal(err)
+               }
+
+               fmt.Printf("Received message  msgId: %s -- content: '%s'\n",
+                       msg.ID(), string(msg.Payload()))
+
+               consumer.Ack(msg)
+       }
+}
diff --git a/pulsar-client-go/examples/producer/producer.go 
b/pulsar-client-go/examples/producer/producer.go
new file mode 100644
index 0000000000..a2505dba92
--- /dev/null
+++ b/pulsar-client-go/examples/producer/producer.go
@@ -0,0 +1,59 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package main
+
+import (
+       "github.com/apache/incubator-pulsar/pulsar-client-go/pulsar"
+       "fmt"
+       "log"
+       "context"
+)
+
+func main() {
+       client, err := pulsar.NewClient(pulsar.ClientOptions{
+               URL:       "pulsar://localhost:6650",
+               IOThreads: 5,
+       })
+
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer client.Close()
+
+       producer, err := client.CreateProducer(pulsar.ProducerOptions{
+               Topic: "my-topic",
+       })
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer producer.Close()
+
+       ctx := context.Background()
+
+       for i := 0; i < 10; i++ {
+               if err := producer.Send(ctx, pulsar.ProducerMessage{
+                       Payload: []byte(fmt.Sprintf("hello-%d", i)),
+               }); err != nil {
+                       log.Fatal(err)
+               }
+       }
+}
diff --git a/pulsar-client-go/examples/reader/reader.go 
b/pulsar-client-go/examples/reader/reader.go
new file mode 100644
index 0000000000..7e76820da1
--- /dev/null
+++ b/pulsar-client-go/examples/reader/reader.go
@@ -0,0 +1,56 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package main
+
+import (
+       "github.com/apache/incubator-pulsar/pulsar-client-go/pulsar"
+       "fmt"
+       "log"
+       "context"
+)
+
+func main() {
+       client, err := pulsar.NewClient(pulsar.ClientOptions{URL: 
"pulsar://localhost:6650"})
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer client.Close()
+
+       reader, err := client.CreateReader(pulsar.ReaderOptions{
+               Topic:          "my-topic",
+               StartMessageID: pulsar.EarliestMessage,
+       })
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       defer reader.Close()
+
+       for {
+               msg, err := reader.Next(context.Background())
+               if err != nil {
+                       log.Fatal(err)
+               }
+
+               fmt.Printf("Received message  msgId: %s -- content: '%s'\n",
+                       msg.ID(), string(msg.Payload()))
+       }
+}
diff --git a/pulsar-client-go/pulsar/c_client.go 
b/pulsar-client-go/pulsar/c_client.go
new file mode 100644
index 0000000000..f2dce18656
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_client.go
@@ -0,0 +1,182 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+/*
+#cgo LDFLAGS: -lpulsar
+#include "c_go_pulsar.h"
+*/
+import "C"
+import (
+       "runtime"
+       "unsafe"
+       "log"
+)
+
+//export pulsarClientLoggerProxy
+func pulsarClientLoggerProxy(level C.pulsar_logger_level_t, file *C.char, line 
C.int, message *C.char, ctx unsafe.Pointer) {
+       logger := restorePointerNoDelete(ctx).(func(LoggerLevel, string, int, 
string))
+
+       logger(LoggerLevel(level), C.GoString(file), int(line), 
C.GoString(message))
+}
+
+func newClient(options ClientOptions) (Client, error) {
+       if options.URL == "" {
+               return nil, newError(C.pulsar_result_InvalidConfiguration, "URL 
is required for client")
+       }
+
+       // Configure the client
+       conf := C.pulsar_client_configuration_create()
+       if options.IOThreads != 0 {
+               C.pulsar_client_configuration_set_io_threads(conf, 
C.int(options.IOThreads))
+       }
+
+       if options.OperationTimeoutSeconds != 0 {
+               
C.pulsar_client_configuration_set_operation_timeout_seconds(conf, 
C.int(options.OperationTimeoutSeconds))
+       }
+
+       if options.MessageListenerThreads != 0 {
+               
C.pulsar_client_configuration_set_message_listener_threads(conf, 
C.int(options.MessageListenerThreads))
+       }
+
+       if options.ConcurrentLookupRequests != 0 {
+               
C.pulsar_client_configuration_set_concurrent_lookup_request(conf, 
C.int(options.ConcurrentLookupRequests))
+       }
+
+       if options.Logger == nil {
+               // Configure a default logger with same date format as Go logs
+               options.Logger = func(level LoggerLevel, file string, line int, 
message string) {
+                       log.Printf("%-5s | %s:%d | %s", level, file, line, 
message)
+               }
+       }
+
+       C._pulsar_client_configuration_set_logger(conf, 
savePointer(options.Logger))
+
+       if options.EnableTLS {
+               C.pulsar_client_configuration_set_use_tls(conf, 
cBool(options.EnableTLS))
+       }
+
+       if options.TLSTrustCertsFilePath != "" {
+               str := C.CString(options.TLSTrustCertsFilePath)
+               defer C.free(unsafe.Pointer(str))
+               
C.pulsar_client_configuration_set_tls_trust_certs_file_path(conf, str)
+       }
+
+       if options.TLSAllowInsecureConnection {
+               
C.pulsar_client_configuration_set_tls_allow_insecure_connection(conf, 
cBool(options.TLSAllowInsecureConnection))
+       }
+
+       if options.StatsIntervalInSeconds != 0 {
+               
C.pulsar_client_configuration_set_stats_interval_in_seconds(conf, 
C.uint(options.StatsIntervalInSeconds))
+       }
+
+       client := &client{
+               ptr: C.pulsar_client_create(C.CString(options.URL), conf),
+       }
+
+       C.pulsar_client_configuration_free(conf)
+       runtime.SetFinalizer(client, clientFinalizer)
+       return client, nil
+}
+
+type client struct {
+       ptr *C.pulsar_client_t
+}
+
+func clientFinalizer(client *client) {
+       C.pulsar_client_free(client.ptr)
+}
+
+func (client *client) CreateProducer(options ProducerOptions) (Producer, 
error) {
+       // Create is implemented on async create with a channel to wait for
+       // completion without blocking the real thread
+       c := make(chan struct {
+               Producer
+               error
+       })
+
+       client.CreateProducerAsync(options, func(producer Producer, err error) {
+               c <- struct {
+                       Producer
+                       error
+               }{producer, err}
+               close(c)
+       })
+
+       res := <-c
+       return res.Producer, res.error
+}
+
+func (client *client) CreateProducerAsync(options ProducerOptions, callback 
func(producer Producer, err error)) {
+       createProducerAsync(client, options, callback)
+}
+
+func (client *client) Subscribe(options ConsumerOptions) (Consumer, error) {
+       c := make(chan struct {
+               Consumer
+               error
+       })
+
+       client.SubscribeAsync(options, func(consumer Consumer, err error) {
+               c <- struct {
+                       Consumer
+                       error
+               }{consumer, err}
+               close(c)
+       })
+
+       res := <-c
+       return res.Consumer, res.error
+}
+
+func (client *client) SubscribeAsync(options ConsumerOptions, callback 
func(Consumer, error)) {
+       subscribeAsync(client, options, callback)
+}
+
+func (client *client) CreateReader(options ReaderOptions) (Reader, error) {
+       c := make(chan struct {
+               Reader
+               error
+       })
+
+       client.CreateReaderAsync(options, func(reader Reader, err error) {
+               c <- struct {
+                       Reader
+                       error
+               }{reader, err}
+               close(c)
+       })
+
+       res := <-c
+       return res.Reader, res.error
+}
+
+func (client *client) CreateReaderAsync(options ReaderOptions, callback 
func(Reader, error)) {
+       createReaderAsync(client, options, callback)
+}
+
+func (client *client) Close() error {
+       res := C.pulsar_client_close(client.ptr)
+       if res != C.pulsar_result_Ok {
+               return newError(res, "Failed to close Pulsar client")
+       } else {
+               return nil
+       }
+}
diff --git a/pulsar-client-go/pulsar/c_consumer.go 
b/pulsar-client-go/pulsar/c_consumer.go
new file mode 100644
index 0000000000..abbe0f194d
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_consumer.go
@@ -0,0 +1,239 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+/*
+#include "c_go_pulsar.h"
+*/
+import "C"
+
+import (
+       "runtime"
+       "time"
+       "unsafe"
+       "context"
+)
+
+type consumer struct {
+       ptr            *C.pulsar_consumer_t
+       defaultChannel chan ConsumerMessage
+}
+
+func consumerFinalizer(c *consumer) {
+       if c.ptr != nil {
+               C.pulsar_consumer_free(c.ptr)
+       }
+}
+
+//export pulsarSubscribeCallbackProxy
+func pulsarSubscribeCallbackProxy(res C.pulsar_result, ptr 
*C.pulsar_consumer_t, ctx unsafe.Pointer) {
+       cc := restorePointer(ctx).(*subscribeContext)
+
+       C.pulsar_consumer_configuration_free(cc.conf)
+
+       if res != C.pulsar_result_Ok {
+               cc.callback(nil, newError(res, "Failed to subscribe to topic"))
+       } else {
+               cc.consumer.ptr = ptr
+               runtime.SetFinalizer(cc.consumer, consumerFinalizer)
+               cc.callback(cc.consumer, nil)
+       }
+}
+
+type subscribeContext struct {
+       conf     *C.pulsar_consumer_configuration_t
+       consumer *consumer
+       callback func(Consumer, error)
+}
+
+func subscribeAsync(client *client, options ConsumerOptions, callback 
func(Consumer, error)) {
+       if options.Topic == "" {
+               go callback(nil, newError(C.pulsar_result_InvalidConfiguration, 
"topic is required"))
+               return
+       }
+
+       if options.SubscriptionName == "" {
+               go callback(nil, newError(C.pulsar_result_InvalidConfiguration, 
"subscription name is required"))
+               return
+       }
+
+       conf := C.pulsar_consumer_configuration_create()
+
+       consumer := &consumer{}
+
+       if options.MessageChannel == nil {
+               // If there is no message listener, set a default channel so 
that we can have receive to
+               // use that
+               consumer.defaultChannel = make(chan ConsumerMessage)
+               options.MessageChannel = consumer.defaultChannel
+       }
+
+       C._pulsar_consumer_configuration_set_message_listener(conf, 
savePointer(&consumerCallback{
+               consumer: consumer,
+               channel:  options.MessageChannel,
+       }))
+
+       if options.AckTimeout != 0 {
+               timeoutMillis := options.AckTimeout.Nanoseconds() / 
int64(time.Millisecond)
+               C.pulsar_consumer_set_unacked_messages_timeout_ms(conf, 
C.ulonglong(timeoutMillis))
+       }
+
+       if options.Type != Exclusive {
+               C.pulsar_consumer_configuration_set_consumer_type(conf, 
C.pulsar_consumer_type(options.Type))
+       }
+
+       // ReceiverQueueSize==0 means to use the default queue size
+       // -1 means to disable the consumer prefetching
+       if options.ReceiverQueueSize > 0 {
+               C.pulsar_consumer_configuration_set_receiver_queue_size(conf, 
C.int(options.ReceiverQueueSize))
+       } else if options.ReceiverQueueSize < 0 {
+               // In C++ client lib, 0 means disable prefetching
+               C.pulsar_consumer_configuration_set_receiver_queue_size(conf, 
C.int(0))
+       }
+
+       if options.MaxTotalReceiverQueueSizeAcrossPartitions != 0 {
+               
C.pulsar_consumer_set_max_total_receiver_queue_size_across_partitions(conf,
+                       
C.int(options.MaxTotalReceiverQueueSizeAcrossPartitions))
+       }
+
+       if options.Name != "" {
+               name := C.CString(options.Name)
+               defer C.free(unsafe.Pointer(name))
+
+               C.pulsar_consumer_set_consumer_name(conf, name)
+       }
+
+       topic := C.CString(options.Topic)
+       subName := C.CString(options.SubscriptionName)
+       defer C.free(unsafe.Pointer(topic))
+       defer C.free(unsafe.Pointer(subName))
+       C._pulsar_client_subscribe_async(client.ptr, topic, subName,
+               conf, savePointer(&subscribeContext{conf: conf, consumer: 
consumer, callback: callback}))
+}
+
+type consumerCallback struct {
+       consumer Consumer
+       channel  chan ConsumerMessage
+}
+
+//export pulsarMessageListenerProxy
+func pulsarMessageListenerProxy(cConsumer *C.pulsar_consumer_t, message 
*C.pulsar_message_t, ctx unsafe.Pointer) {
+       cc := restorePointerNoDelete(ctx).(*consumerCallback)
+
+       defer func() {
+               ex := recover()
+               if ex != nil {
+                       // There was an error when sending channel (eg: already 
closed)
+               }
+       }()
+
+       cc.channel <- ConsumerMessage{cc.consumer, newMessageWrapper(message)}
+}
+
+//// Consumer
+
+func (c *consumer) Topic() string {
+       return C.GoString(C.pulsar_consumer_get_topic(c.ptr))
+}
+
+func (c *consumer) Subscription() string {
+       return C.GoString(C.pulsar_consumer_get_subscription_name(c.ptr))
+}
+
+func (c *consumer) Unsubscribe() error {
+       channel := make(chan error)
+       c.UnsubscribeAsync(func(err error) {
+               channel <- err; close(channel) })
+       return <-channel
+}
+
+func (c *consumer) UnsubscribeAsync(callback func(error)) {
+       C._pulsar_consumer_unsubscribe_async(c.ptr, savePointer(callback))
+}
+
+//export pulsarConsumerUnsubscribeCallbackProxy
+func pulsarConsumerUnsubscribeCallbackProxy(res C.pulsar_result, ctx 
unsafe.Pointer) {
+       callback := restorePointer(ctx).(func(err error))
+
+       if res != C.pulsar_result_Ok {
+               go callback(newError(res, "Failed to unsubscribe consumer"))
+       } else {
+               go callback(nil)
+       }
+}
+
+func (c *consumer) Receive(ctx context.Context) (Message, error) {
+       select {
+       case <-ctx.Done():
+               return nil, ctx.Err()
+
+       case cm := <-c.defaultChannel:
+               return cm.Message, nil
+       }
+}
+
+func (c *consumer) Ack(msg Message) error {
+       C.pulsar_consumer_acknowledge_async(c.ptr, msg.(*message).ptr, nil, nil)
+       return nil
+}
+
+func (c *consumer) AckID(msgId MessageID) error {
+       C.pulsar_consumer_acknowledge_async_id(c.ptr, msgId.(*messageID).ptr, 
nil, nil)
+       return nil
+}
+
+func (c *consumer) AckCumulative(msg Message) error {
+       C.pulsar_consumer_acknowledge_cumulative_async(c.ptr, 
msg.(*message).ptr, nil, nil)
+       return nil
+}
+
+func (c *consumer) AckCumulativeID(msgId MessageID) error {
+       C.pulsar_consumer_acknowledge_cumulative_async_id(c.ptr, 
msgId.(*messageID).ptr, nil, nil)
+       return nil
+}
+
+func (c *consumer) Close() error {
+       channel := make(chan error)
+       c.CloseAsync(func(err error) { channel <- err; close(channel) })
+       return <-channel
+}
+
+func (c *consumer) CloseAsync(callback func(error)) {
+       if c.defaultChannel != nil {
+               close(c.defaultChannel)
+       }
+
+       C._pulsar_consumer_close_async(c.ptr, savePointer(callback))
+}
+
+//export pulsarConsumerCloseCallbackProxy
+func pulsarConsumerCloseCallbackProxy(res C.pulsar_result, ctx unsafe.Pointer) 
{
+       callback := restorePointer(ctx).(func(err error))
+
+       if res != C.pulsar_result_Ok {
+               go callback(newError(res, "Failed to close Consumer"))
+       } else {
+               go callback(nil)
+       }
+}
+
+func (c *consumer) RedeliverUnackedMessages() {
+       C.pulsar_consumer_redeliver_unacknowledged_messages(c.ptr)
+}
diff --git a/pulsar-client-go/pulsar/c_error.go 
b/pulsar-client-go/pulsar/c_error.go
new file mode 100644
index 0000000000..a7c83ae518
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_error.go
@@ -0,0 +1,60 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import "C"
+
+/*
+#include "c_go_pulsar.h"
+*/
+import "C"
+import "fmt"
+
+func cBool(flag bool) C.int {
+       if flag {
+               return C.int(1)
+       } else {
+               return C.int(0)
+       }
+}
+
+type Error struct {
+       msg    string
+       result Result
+}
+
+func newError(result C.pulsar_result, msg string) error {
+       return &Error{
+               msg:    fmt.Sprintf("%s: %s", msg, 
C.GoString(C.pulsar_result_str(result))),
+               result: Result(result),
+       }
+}
+
+func (e *Error) Result() Result {
+       return e.result
+}
+
+func (e *Error) Error() string {
+       return e.msg
+}
+
+func (r Result) String() string {
+       return C.GoString(C.pulsar_result_str(C.pulsar_result(r)))
+}
diff --git a/pulsar-client-go/pulsar/c_go_pulsar.h 
b/pulsar-client-go/pulsar/c_go_pulsar.h
new file mode 100644
index 0000000000..d2ca7ba632
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_go_pulsar.h
@@ -0,0 +1,136 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+#pragma once
+
+#include <pulsar/c/client.h>
+#include <stdlib.h>
+
+// Callback proxy functions
+
+void pulsarClientLoggerProxy(pulsar_logger_level_t level, char* file, int 
line, char* message, void *ctx);
+
+static inline void pulsarClientLoggerConstProxy(pulsar_logger_level_t level, 
const char* file, int line, const char* message, void *ctx) {
+    pulsarClientLoggerProxy(level, (char*)file, line, (char*)message, ctx);
+}
+
+static inline void 
_pulsar_client_configuration_set_logger(pulsar_client_configuration_t *conf, 
void *ctx) {
+    pulsar_client_configuration_set_logger(conf, pulsarClientLoggerConstProxy, 
ctx);
+}
+
+void pulsarCreateProducerCallbackProxy(pulsar_result result, pulsar_producer_t 
*producer, void *ctx);
+
+static inline void _pulsar_client_create_producer_async(pulsar_client_t 
*client, const char *topic,
+                                                        const 
pulsar_producer_configuration_t *conf,
+                                                        void *ctx) {
+    pulsar_client_create_producer_async(client, topic, conf, 
pulsarCreateProducerCallbackProxy, ctx);
+}
+
+void pulsarProducerCloseCallbackProxy(pulsar_result result, void *ctx);
+
+static inline void _pulsar_producer_close_async(pulsar_producer_t *producer, 
void *ctx) {
+    pulsar_producer_close_async(producer, pulsarProducerCloseCallbackProxy, 
ctx);
+}
+
+void pulsarProducerSendCallbackProxy(pulsar_result result, pulsar_message_t 
*message, void *ctx);
+
+static inline void _pulsar_producer_send_async(pulsar_producer_t *producer, 
pulsar_message_t *message,
+                                               void *ctx) {
+    pulsar_producer_send_async(producer, message, 
pulsarProducerSendCallbackProxy, ctx);
+}
+
+int pulsarRouterCallbackProxy(pulsar_message_t *msg, pulsar_topic_metadata_t 
*topicMetadata, void* ctx);
+
+
+static inline void 
_pulsar_producer_configuration_set_message_router(pulsar_producer_configuration_t
 *conf, void *ctx) {
+    pulsar_producer_configuration_set_message_router(conf, 
pulsarRouterCallbackProxy, ctx);
+}
+
+//// Consumer callbacks
+
+void pulsarSubscribeCallbackProxy(pulsar_result result, pulsar_consumer_t 
*consumer, void *ctx);
+
+static inline void _pulsar_client_subscribe_async(pulsar_client_t *client, 
const char *topic,
+                                                  const char *subscriptionName,
+                                                  const 
pulsar_consumer_configuration_t *conf, void *ctx) {
+    pulsar_client_subscribe_async(client, topic, subscriptionName, conf, 
pulsarSubscribeCallbackProxy, ctx);
+}
+
+void pulsarMessageListenerProxy(pulsar_consumer_t *consumer, pulsar_message_t 
*message, void *ctx);
+
+static inline void _pulsar_consumer_configuration_set_message_listener(
+    pulsar_consumer_configuration_t *consumer_configuration, void *ctx) {
+    pulsar_consumer_configuration_set_message_listener(consumer_configuration, 
pulsarMessageListenerProxy,
+                                                       ctx);
+}
+
+void pulsarConsumerUnsubscribeCallbackProxy(pulsar_result result, void *ctx);
+
+static inline void _pulsar_consumer_unsubscribe_async(pulsar_consumer_t 
*consumer, void *ctx) {
+    pulsar_consumer_unsubscribe_async(consumer, 
pulsarConsumerUnsubscribeCallbackProxy, ctx);
+}
+
+void pulsarConsumerCloseCallbackProxy(pulsar_result result, void *ctx);
+
+static inline void _pulsar_consumer_close_async(pulsar_consumer_t *consumer, 
void *ctx) {
+    pulsar_consumer_close_async(consumer, pulsarConsumerCloseCallbackProxy, 
ctx);
+}
+
+//// Reader callbacks
+
+void pulsarCreateReaderCallbackProxy(pulsar_result result, pulsar_reader_t 
*reader, void *ctx);
+
+static inline void _pulsar_client_create_reader_async(pulsar_client_t *client, 
const char *topic,
+                                                      const 
pulsar_message_id_t *startMessageId,
+                                                      
pulsar_reader_configuration_t *conf, void *ctx) {
+    pulsar_client_create_reader_async(client, topic, startMessageId, conf, 
pulsarCreateReaderCallbackProxy,
+                                      ctx);
+}
+
+void pulsarReaderListenerProxy(pulsar_reader_t *reader, pulsar_message_t 
*message, void *ctx);
+
+static inline void _pulsar_reader_configuration_set_reader_listener(
+    pulsar_reader_configuration_t *reader_configuration, void *ctx) {
+    pulsar_reader_configuration_set_reader_listener(reader_configuration, 
pulsarReaderListenerProxy, ctx);
+}
+
+void pulsarReaderCloseCallbackProxy(pulsar_result result, void *ctx);
+
+static inline void _pulsar_reader_close_async(pulsar_reader_t *reader, void 
*ctx) {
+    pulsar_reader_close_async(reader, pulsarReaderCloseCallbackProxy, ctx);
+}
+
+
+//// String array manipulation
+
+static char** newStringArray(int size) {
+    return calloc(sizeof(char*), size);
+}
+
+static void setString(char** array, char *str, int n) {
+    array[n] = str;
+}
+
+static void freeStringArray(char* *array, int size) {
+    for (int i = 0; i < size; i++) {
+        free(array[i]);
+    }
+
+    free(array);
+}
diff --git a/pulsar-client-go/pulsar/c_message.go 
b/pulsar-client-go/pulsar/c_message.go
new file mode 100644
index 0000000000..4f391ad6ac
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_message.go
@@ -0,0 +1,207 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+/*
+#include "c_go_pulsar.h"
+*/
+import "C"
+
+import (
+       "reflect"
+       "runtime"
+       "unsafe"
+       "time"
+)
+
+type message struct {
+       ptr *C.pulsar_message_t
+}
+
+type messageID struct {
+       ptr *C.pulsar_message_id_t
+}
+
+////////////////////////////////////////////////////////////
+
+func buildMessage(message ProducerMessage) *C.pulsar_message_t {
+
+       cMsg := C.pulsar_message_create()
+
+       if message.Key != "" {
+               cKey := C.CString(message.Key)
+               defer C.free(unsafe.Pointer(cKey))
+               C.pulsar_message_set_partition_key(cMsg, cKey)
+       }
+
+       if message.Payload != nil {
+               C.pulsar_message_set_content(cMsg, 
unsafe.Pointer(&message.Payload[0]), C.ulong(len(message.Payload)))
+       }
+
+       if message.Properties != nil {
+               for key, value := range message.Properties {
+                       cKey := C.CString(key)
+                       cValue := C.CString(value)
+
+                       C.pulsar_message_set_property(cMsg, cKey, cValue)
+
+                       C.free(unsafe.Pointer(cKey))
+                       C.free(unsafe.Pointer(cValue))
+               }
+       }
+
+       if message.EventTime.UnixNano() != 0 {
+               C.pulsar_message_set_event_timestamp(cMsg, 
timeToUnixTimestampMillis(message.EventTime))
+       }
+
+       if message.ReplicationClusters != nil {
+               if len(message.ReplicationClusters) == 0 {
+                       // Empty list means to disable replication
+                       C.pulsar_message_disable_replication(cMsg, C.int(1))
+               } else {
+                       size := C.int(len(message.ReplicationClusters))
+                       array := C.newStringArray(size)
+                       defer C.freeStringArray(array, size)
+
+                       for i, s := range message.ReplicationClusters {
+                               C.setString(array, C.CString(s), C.int(i))
+                       }
+
+                       C.pulsar_message_set_replication_clusters(cMsg, array)
+               }
+       }
+
+       return cMsg
+}
+
+////////////// Message
+
+func newMessageWrapper(ptr *C.pulsar_message_t) Message {
+       msg := &message{ptr: ptr}
+       runtime.SetFinalizer(msg, messageFinalizer)
+       return msg
+}
+
+func messageFinalizer(msg *message) {
+       C.pulsar_message_free(msg.ptr)
+}
+
+func (m *message) Properties() map[string]string {
+       cProperties := C.pulsar_message_get_properties(m.ptr)
+       defer C.pulsar_string_map_free(cProperties)
+
+       properties := make(map[string]string)
+       count := int(C.pulsar_string_map_size(cProperties))
+       for i := 0; i < count; i++ {
+               key := C.GoString(C.pulsar_string_map_get_key(cProperties, 
C.int(i)))
+               value := C.GoString(C.pulsar_string_map_get_value(cProperties, 
C.int(i)))
+
+               properties[key] = value
+       }
+
+       return properties
+}
+
+func (m *message) Payload() []byte {
+       payload := C.pulsar_message_get_data(m.ptr)
+       size := C.pulsar_message_get_length(m.ptr)
+
+       // Get the byte array without copying the data. The array will be valid
+       // until we free the message in m.ptr
+       slice := &reflect.SliceHeader{Data: uintptr(payload), Len: int(size), 
Cap: int(size)}
+       return *(*[]byte)(unsafe.Pointer(slice))
+}
+
+func (m *message) ID() MessageID {
+       return newMessageId(m.ptr)
+}
+
+func (m *message) PublishTime() time.Time {
+       return 
timeFromUnixTimestampMillis(C.pulsar_message_get_publish_timestamp(m.ptr))
+}
+
+func (m *message) EventTime() *time.Time {
+       eventTime := C.pulsar_message_get_event_timestamp(m.ptr)
+       if uint64(eventTime) == 0 {
+               return nil
+       } else {
+               res := timeFromUnixTimestampMillis(eventTime)
+               return &res
+       }
+}
+
+func (m *message) Key() string {
+       return C.GoString(C.pulsar_message_get_partitionKey(m.ptr))
+}
+
+//////// MessageID
+
+func newMessageId(msg *C.pulsar_message_t) MessageID {
+       msgId := &messageID{ptr: C.pulsar_message_get_message_id(msg)}
+       runtime.SetFinalizer(msgId, messageIdFinalizer)
+       return msgId
+}
+
+func messageIdFinalizer(msgID *messageID) {
+       C.pulsar_message_id_free(msgID.ptr)
+}
+
+func (m *messageID) Serialize() []byte {
+       var size C.int
+       buf := C.pulsar_message_id_serialize(m.ptr, &size)
+       defer C.free(unsafe.Pointer(buf))
+       return C.GoBytes(buf, size)
+}
+
+func deserializeMessageId(data []byte) MessageID {
+       msgId := &messageID{ptr: 
C.pulsar_message_id_deserialize(unsafe.Pointer(&data[0]), C.uint(len(data)))}
+       runtime.SetFinalizer(msgId, messageIdFinalizer)
+       return msgId
+}
+
+func (m *messageID) String() string {
+       str := C.pulsar_message_id_str(m.ptr)
+       defer C.free(unsafe.Pointer(str))
+       return C.GoString(str)
+}
+
+func earliestMessageID() *messageID {
+       // No need to use finalizer since the pointer doesn't need to be freed
+       return &messageID{C.pulsar_message_id_earliest()}
+}
+
+func latestMessageID() *messageID {
+       // No need to use finalizer since the pointer doesn't need to be freed
+       return &messageID{C.pulsar_message_id_latest()}
+}
+
+func timeFromUnixTimestampMillis(timestamp C.ulonglong) time.Time {
+       ts := int64(timestamp)
+       seconds := ts / int64(time.Millisecond)
+       millis := ts - seconds
+       nanos := millis * int64(time.Millisecond)
+       return time.Unix(seconds, nanos)
+}
+
+func timeToUnixTimestampMillis(t time.Time) C.ulonglong {
+       nanos := t.UnixNano()
+       millis := nanos / int64(time.Millisecond)
+       return C.ulonglong(millis)
+}
diff --git a/pulsar-client-go/pulsar/c_producer.go 
b/pulsar-client-go/pulsar/c_producer.go
new file mode 100644
index 0000000000..b4cd2c5e14
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_producer.go
@@ -0,0 +1,222 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+/*
+#include "c_go_pulsar.h"
+*/
+import "C"
+import (
+       "runtime"
+       "unsafe"
+       "time"
+       "context"
+)
+
+type createProducerCtx struct {
+       callback func(producer Producer, err error)
+       conf     *C.pulsar_producer_configuration_t
+}
+
+//export pulsarCreateProducerCallbackProxy
+func pulsarCreateProducerCallbackProxy(res C.pulsar_result, ptr 
*C.pulsar_producer_t, ctx unsafe.Pointer) {
+       producerCtx := restorePointer(ctx).(createProducerCtx)
+
+       C.pulsar_producer_configuration_free(producerCtx.conf)
+
+       if res != C.pulsar_result_Ok {
+               producerCtx.callback(nil, newError(res, "Failed to create 
Producer"))
+       } else {
+               p := &producer{ptr: ptr}
+               runtime.SetFinalizer(p, producerFinalizer)
+               producerCtx.callback(p, nil)
+       }
+}
+
+func createProducerAsync(client *client, options ProducerOptions, callback 
func(producer Producer, err error)) {
+       if options.Topic == "" {
+               go callback(nil, newError(C.pulsar_result_InvalidConfiguration, 
"topic is required when creating producer"))
+               return
+       }
+
+       conf := C.pulsar_producer_configuration_create()
+
+       if options.Name != "" {
+               cName := C.CString(options.Name)
+               defer C.free(unsafe.Pointer(cName))
+               C.pulsar_producer_configuration_set_producer_name(conf, cName)
+       }
+
+       // If SendTimeout is 0, we'll leave the default configured value on C 
library
+       if options.SendTimeout > 0 {
+               timeoutMillis := options.SendTimeout.Nanoseconds() / 
int64(time.Millisecond)
+               C.pulsar_producer_configuration_set_send_timeout(conf, 
C.int(timeoutMillis))
+       } else if options.SendTimeout < 0 {
+               // Set infinite publish timeout, which is specified as 0 in C 
API
+               C.pulsar_producer_configuration_set_send_timeout(conf, C.int(0))
+       }
+
+       if options.MaxPendingMessages != 0 {
+               C.pulsar_producer_configuration_set_max_pending_messages(conf, 
C.int(options.MaxPendingMessages))
+       }
+
+       if options.MaxPendingMessagesAcrossPartitions != 0 {
+               
C.pulsar_producer_configuration_set_max_pending_messages_across_partitions(conf,
 C.int(options.MaxPendingMessagesAcrossPartitions))
+       }
+
+       if options.BlockIfQueueFull {
+               C.pulsar_producer_configuration_set_block_if_queue_full(conf, 
cBool(options.BlockIfQueueFull))
+       }
+
+       switch options.MessageRoutingMode {
+       case RoundRobinDistribution:
+               
C.pulsar_producer_configuration_set_partitions_routing_mode(conf, 
C.pulsar_RoundRobinDistribution)
+       case UseSinglePartition:
+               
C.pulsar_producer_configuration_set_partitions_routing_mode(conf, 
C.pulsar_UseSinglePartition)
+       case CustomPartition:
+               
C.pulsar_producer_configuration_set_partitions_routing_mode(conf, 
C.pulsar_CustomPartition)
+       }
+
+       switch options.HashingScheme {
+       case JavaStringHash:
+               C.pulsar_producer_configuration_set_hashing_scheme(conf, 
C.pulsar_JavaStringHash)
+       case Murmur3_32Hash:
+               C.pulsar_producer_configuration_set_hashing_scheme(conf, 
C.pulsar_Murmur3_32Hash)
+       case BoostHash:
+               C.pulsar_producer_configuration_set_hashing_scheme(conf, 
C.pulsar_BoostHash)
+       }
+
+       if options.CompressionType != NoCompression {
+               C.pulsar_producer_configuration_set_compression_type(conf, 
C.pulsar_compression_type(options.CompressionType))
+       }
+
+       if options.MessageRouter != nil {
+               C._pulsar_producer_configuration_set_message_router(conf, 
savePointer(&options.MessageRouter))
+       }
+
+       if options.Batching {
+               C.pulsar_producer_configuration_set_batching_enabled(conf, 
cBool(options.Batching))
+       }
+
+       if options.BatchingMaxPublishDelay != 0 {
+               delayMillis := options.BatchingMaxPublishDelay.Nanoseconds() / 
int64(time.Millisecond)
+               
C.pulsar_producer_configuration_set_batching_max_publish_delay_ms(conf, 
C.ulong(delayMillis))
+       }
+
+       if options.BatchingMaxMessages != 0 {
+               C.pulsar_producer_configuration_set_batching_max_messages(conf, 
C.uint(options.BatchingMaxMessages))
+       }
+
+       topicName := C.CString(options.Topic)
+       defer C.free(unsafe.Pointer(topicName))
+
+       C._pulsar_client_create_producer_async(client.ptr, topicName, conf,
+               savePointer(createProducerCtx{callback, conf}))
+}
+
+type topicMetadata struct {
+       numPartitions int
+}
+
+func (tm *topicMetadata) NumPartitions() int {
+       return tm.numPartitions
+}
+
+//export pulsarRouterCallbackProxy
+func pulsarRouterCallbackProxy(msg *C.pulsar_message_t, metadata 
*C.pulsar_topic_metadata_t, ctx unsafe.Pointer) C.int {
+       router := restorePointerNoDelete(ctx).(*func(msg Message, metadata 
TopicMetadata) int)
+       partitionIdx := (*router)(&message{msg}, 
&topicMetadata{int(C.pulsar_topic_metadata_get_num_partitions(metadata))})
+       return C.int(partitionIdx)
+}
+
+/// Producer
+
+type producer struct {
+       ptr *C.pulsar_producer_t
+}
+
+func producerFinalizer(p *producer) {
+       C.pulsar_producer_free(p.ptr)
+}
+
+func (p *producer) Topic() string {
+       return C.GoString(C.pulsar_producer_get_topic(p.ptr))
+}
+
+func (p *producer) Name() string {
+       return C.GoString(C.pulsar_producer_get_producer_name(p.ptr))
+}
+
+func (p *producer) Send(ctx context.Context, msg ProducerMessage) error {
+       c := make(chan error)
+       p.SendAsync(ctx, msg, func(msg ProducerMessage, err error) { c <- err; 
close(c) })
+
+       select {
+       case <-ctx.Done():
+               return ctx.Err()
+
+       case cm := <-c:
+               return cm
+       }
+}
+
+type sendCallback struct {
+       message ProducerMessage
+       callback func(ProducerMessage, error)
+}
+
+//export pulsarProducerSendCallbackProxy
+func pulsarProducerSendCallbackProxy(res C.pulsar_result, message 
*C.pulsar_message_t, ctx unsafe.Pointer) {
+       sendCallback := restorePointer(ctx).(sendCallback)
+
+       if res != C.pulsar_result_Ok {
+               sendCallback.callback(sendCallback.message, newError(res, 
"Failed to send message"))
+       } else {
+               sendCallback.callback(sendCallback.message, nil)
+       }
+}
+
+func (p *producer) SendAsync(ctx context.Context, msg ProducerMessage, 
callback func(ProducerMessage, error)) {
+       cMsg := buildMessage(msg)
+       defer C.pulsar_message_free(cMsg)
+
+       C._pulsar_producer_send_async(p.ptr, cMsg, 
savePointer(sendCallback{msg, callback}))
+}
+
+func (p *producer) Close() error {
+       c := make(chan error)
+       p.CloseAsync(func(err error) { c <- err; close(c) })
+       return <-c
+}
+
+func (p *producer) CloseAsync(callback func(error)) {
+       C._pulsar_producer_close_async(p.ptr, savePointer(callback))
+}
+
+//export pulsarProducerCloseCallbackProxy
+func pulsarProducerCloseCallbackProxy(res C.pulsar_result, ctx unsafe.Pointer) 
{
+       callback := restorePointer(ctx).(func(error))
+
+       if res != C.pulsar_result_Ok {
+               callback(newError(res, "Failed to close Producer"))
+       } else {
+               callback(nil)
+       }
+}
diff --git a/pulsar-client-go/pulsar/c_reader.go 
b/pulsar-client-go/pulsar/c_reader.go
new file mode 100644
index 0000000000..730f9b86fe
--- /dev/null
+++ b/pulsar-client-go/pulsar/c_reader.go
@@ -0,0 +1,172 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+/*
+#include "c_go_pulsar.h"
+*/
+import "C"
+
+import (
+       "runtime"
+       "unsafe"
+       "context"
+)
+
+type reader struct {
+       ptr            *C.pulsar_reader_t
+       defaultChannel chan ReaderMessage
+}
+
+func readerFinalizer(c *reader) {
+       if c.ptr != nil {
+               C.pulsar_reader_free(c.ptr)
+       }
+}
+
+//export pulsarCreateReaderCallbackProxy
+func pulsarCreateReaderCallbackProxy(res C.pulsar_result, ptr 
*C.pulsar_reader_t, ctx unsafe.Pointer) {
+       cc := restorePointer(ctx).(*readerAndCallback)
+
+       C.pulsar_reader_configuration_free(cc.conf)
+
+       if res != C.pulsar_result_Ok {
+               cc.callback(nil, newError(res, "Failed to create Reader"))
+       } else {
+               cc.reader.ptr = ptr
+               runtime.SetFinalizer(cc.reader, readerFinalizer)
+               cc.callback(cc.reader, nil)
+       }
+}
+
+type readerAndCallback struct {
+       reader   *reader
+       conf     *C.pulsar_reader_configuration_t
+       callback func(Reader, error)
+}
+
+func createReaderAsync(client *client, options ReaderOptions, callback 
func(Reader, error)) {
+       if options.Topic == "" {
+               go callback(nil, newError(C.pulsar_result_InvalidConfiguration, 
"topic is required"))
+               return
+       }
+
+       if options.StartMessageID == nil {
+               go callback(nil, newError(C.pulsar_result_InvalidConfiguration, 
"start message id is required"))
+               return
+       }
+
+       reader := &reader{}
+
+       if options.MessageChannel == nil {
+               // If there is no message listener, set a default channel so 
that we can have receive to
+               // use that
+               reader.defaultChannel = make(chan ReaderMessage)
+               options.MessageChannel = reader.defaultChannel
+       }
+
+       conf := C.pulsar_reader_configuration_create()
+
+       C._pulsar_reader_configuration_set_reader_listener(conf, 
savePointer(&readerCallback{
+               reader:  reader,
+               channel: options.MessageChannel,
+       }))
+
+       if options.ReceiverQueueSize != 0 {
+               C.pulsar_reader_configuration_set_receiver_queue_size(conf, 
C.int(options.ReceiverQueueSize))
+       }
+
+       if options.SubscriptionRolePrefix != "" {
+               prefix := C.CString(options.SubscriptionRolePrefix)
+               defer C.free(unsafe.Pointer(prefix))
+               
C.pulsar_reader_configuration_set_subscription_role_prefix(conf, prefix)
+       }
+
+       if options.Name != "" {
+               name := C.CString(options.Name)
+               defer C.free(unsafe.Pointer(name))
+
+               C.pulsar_reader_configuration_set_reader_name(conf, name)
+       }
+
+       topic := C.CString(options.Topic)
+       defer C.free(unsafe.Pointer(topic))
+
+       C._pulsar_client_create_reader_async(client.ptr, topic, 
options.StartMessageID.(*messageID).ptr,
+               conf, savePointer(&readerAndCallback{reader, conf, callback}))
+}
+
+type readerCallback struct {
+       reader  Reader
+       channel chan ReaderMessage
+}
+
+//export pulsarReaderListenerProxy
+func pulsarReaderListenerProxy(cReader *C.pulsar_reader_t, message 
*C.pulsar_message_t, ctx unsafe.Pointer) {
+       rc := restorePointerNoDelete(ctx).(*readerCallback)
+
+       defer func() {
+               ex := recover()
+               if ex != nil {
+                       // There was an error when sending channel (eg: already 
closed)
+               }
+       }()
+
+       rc.channel <- ReaderMessage{rc.reader, newMessageWrapper(message)}
+}
+
+func (r *reader) Topic() string {
+       return C.GoString(C.pulsar_reader_get_topic(r.ptr))
+}
+
+func (r *reader) Next(ctx context.Context) (Message, error) {
+       select {
+       case <-ctx.Done():
+               return nil, ctx.Err()
+
+       case rm := <-r.defaultChannel:
+               return rm.Message, nil
+       }
+}
+
+func (r *reader) Close() error {
+       channel := make(chan error)
+       r.CloseAsync(func(err error) { channel <- err; close(channel) })
+       return <-channel
+}
+
+func (r *reader) CloseAsync(callback func(error)) {
+       if r.defaultChannel != nil {
+               close(r.defaultChannel)
+       }
+
+       C._pulsar_reader_close_async(r.ptr, savePointer(callback))
+}
+
+//export pulsarReaderCloseCallbackProxy
+func pulsarReaderCloseCallbackProxy(res C.pulsar_result, ctx unsafe.Pointer) {
+       callback := restorePointer(ctx).(func(err error))
+
+       if res != C.pulsar_result_Ok {
+               callback(newError(res, "Failed to close Reader"))
+       } else {
+               callback(nil)
+       }
+}
diff --git a/pulsar-client-go/pulsar/client.go 
b/pulsar-client-go/pulsar/client.go
new file mode 100644
index 0000000000..b814392203
--- /dev/null
+++ b/pulsar-client-go/pulsar/client.go
@@ -0,0 +1,87 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import "time"
+
+func NewClient(options ClientOptions) (Client, error) {
+       return newClient(options)
+}
+
+// Builder interface that is used to construct a Pulsar Client instance.
+type ClientOptions struct {
+       // Configure the service URL for the Pulsar service.
+       // This parameter is required
+       URL string
+
+       // Number of threads to be used for handling connections to brokers 
(default: 1 thread)
+       IOThreads int
+
+       // Set the operation timeout (default: 30 seconds)
+       // Producer-create, subscribe and unsubscribe operations will be 
retried until this interval, after which the
+       // operation will be maked as failed
+       OperationTimeoutSeconds time.Duration
+
+       // Set the number of threads to be used for message listeners (default: 
1 thread)
+       MessageListenerThreads int
+
+       // Number of concurrent lookup-requests allowed to send on each 
broker-connection to prevent overload on broker.
+       // (default: 5000) It should be configured with higher value only in 
case of it requires to produce/subscribe
+       // on thousands of topic using created Pulsar Client
+       ConcurrentLookupRequests int
+
+       // Provide a custom logger implementation where all Pulsar library 
info/warn/error messages will be routed
+       // By default, log messages will be printed on standard output. By 
passing a logger function, application
+       // can determine how to print logs. This function will be called each 
time the Pulsar client library wants
+       // to write any logs.
+       Logger func(level LoggerLevel, file string, line int, message string)
+
+       // Configure whether to use TLS encryption on the connection (default: 
false)
+       EnableTLS bool
+
+       // Set the path to the trusted TLS certificate file
+       TLSTrustCertsFilePath string
+
+       // Configure whether the Pulsar client accept untrusted TLS certificate 
from broker (default: false)
+       TLSAllowInsecureConnection bool
+
+       // Set the interval between each stat info (default: 60 seconds). Stats 
will be activated with positive
+       // statsIntervalSeconds It should be set to at least 1 second
+       StatsIntervalInSeconds int
+}
+
+type Client interface {
+       // Create the producer instance
+       // This method will block until the producer is created successfully
+       CreateProducer(ProducerOptions) (Producer, error)
+
+       // Create a `Consumer` by subscribing to a topic.
+       //
+       // If the subscription does not exist, a new subscription will be 
created and all messages published after the
+       // creation will be retained until acknowledged, even if the consumer 
is not connected
+       Subscribe(ConsumerOptions) (Consumer, error)
+
+       // Create a Reader instance.
+       // This method will block until the reader is created successfully.
+       CreateReader(ReaderOptions) (Reader, error)
+
+       // Close the Client and free associated resources
+       Close() error
+}
diff --git a/pulsar-client-go/pulsar/consumer.go 
b/pulsar-client-go/pulsar/consumer.go
new file mode 100644
index 0000000000..4ce0857652
--- /dev/null
+++ b/pulsar-client-go/pulsar/consumer.go
@@ -0,0 +1,136 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import (
+       "time"
+       "context"
+)
+
+// Pair of a Consumer and Message
+type ConsumerMessage struct {
+       Consumer
+       Message
+}
+
+// Types of subscription supported by Pulsar
+type SubscriptionType int
+
+const (
+       // There can be only 1 consumer on the same topic with the same 
subscription name
+       Exclusive SubscriptionType = 0
+
+       // Multiple consumer will be able to use the same subscription name and 
the messages will be dispatched according to
+       // a round-robin rotation between the connected consumers
+       Shared SubscriptionType = 1
+
+       // Multiple consumer will be able to use the same subscription name but 
only 1 consumer will receive the messages.
+       // If that consumer disconnects, one of the other connected consumers 
will start receiving messages.
+       Failover SubscriptionType = 2
+)
+
+// ConsumerBuilder is used to configure and create instances of Consumer
+type ConsumerOptions struct {
+       // Specify the topic this consumer will subscribe on.
+       // This argument is required when subscribing
+       Topic string
+
+       // Specify the subscription name for this consumer
+       // This argument is required when subscribing
+       SubscriptionName string
+
+       // Set the timeout for unacked messages
+       // Message not acknowledged within the give time, will be replayed by 
the broker to the same or a different consumer
+       // Default is 0, which means message are not being replayed based on 
ack time
+       AckTimeout time.Duration
+
+       // Select the subscription type to be used when subscribing to the 
topic.
+       // Default is `Exclusive`
+       Type SubscriptionType
+
+       // Sets a `MessageChannel` for the consumer
+       // When a message is received, it will be pushed to the channel for 
consumption
+       MessageChannel chan ConsumerMessage
+
+       // Sets the size of the consumer receive queue.
+       // The consumer receive queue controls how many messages can be 
accumulated by the `Consumer` before the
+       // application calls `Consumer.receive()`. Using a higher value could 
potentially increase the consumer
+       // throughput at the expense of bigger memory utilization.
+       // Default value is `1000` messages and should be good for most use 
cases.
+       // Set to -1 to disable prefetching in consumer
+       ReceiverQueueSize int
+
+       // Set the max total receiver queue size across partitions.
+       // This setting will be used to reduce the receiver queue size for 
individual partitions
+       // ReceiverQueueSize(int) if the total exceeds this value (default: 
50000).
+       MaxTotalReceiverQueueSizeAcrossPartitions int
+
+       // Set the consumer name.
+       Name string
+}
+
+// An interface that abstracts behavior of Pulsar's consumer
+type Consumer interface {
+       // Get the topic for the consumer
+       Topic() string
+
+       // Get a subscription for the consumer
+       Subscription() string
+
+       // Unsubscribe the consumer
+       Unsubscribe() error
+
+       // Receives a single message.
+       // This calls blocks until a message is available.
+       Receive(context.Context) (Message, error)
+
+       //Ack the consumption of a single message
+       Ack(Message) error
+
+       // Ack the consumption of a single message, identified by its MessageID
+       AckID(MessageID) error
+
+       // Ack the reception of all the messages in the stream up to (and 
including) the provided message.
+       // This method will block until the acknowledge has been sent to the 
broker. After that, the messages will not be
+       // re-delivered to this consumer.
+       //
+       // Cumulative acknowledge cannot be used when the consumer type is set 
to ConsumerShared.
+       //
+       // It's equivalent to calling asyncAcknowledgeCumulative(Message) and 
waiting for the callback to be triggered.
+       AckCumulative(Message) error
+
+       // Ack the reception of all the messages in the stream up to (and 
including) the provided message.
+       // This method will block until the acknowledge has been sent to the 
broker. After that, the messages will not be
+       // re-delivered to this consumer.
+       //
+       // Cumulative acknowledge cannot be used when the consumer type is set 
to ConsumerShared.
+       //
+       // It's equivalent to calling asyncAcknowledgeCumulative(MessageID) and 
waiting for the callback to be triggered.
+       AckCumulativeID(MessageID) error
+
+       // Close the consumer and stop the broker to push more messages
+       Close() error
+
+       // Redelivers all the unacknowledged messages. In Failover mode, the 
request is ignored if the consumer is not
+       // active for the given topic. In Shared mode, the consumers messages 
to be redelivered are distributed across all
+       // the connected consumers. This is a non blocking call and doesn't 
throw an exception. In case the connection
+       // breaks, the messages are redelivered after reconnect.
+       RedeliverUnackedMessages()
+}
diff --git a/pulsar-client-go/pulsar/consumer_test.go 
b/pulsar-client-go/pulsar/consumer_test.go
new file mode 100644
index 0000000000..2930f19d9c
--- /dev/null
+++ b/pulsar-client-go/pulsar/consumer_test.go
@@ -0,0 +1,133 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import (
+       "testing"
+       "fmt"
+       "context"
+       "time"
+)
+
+func TestConsumerConnectError(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://invalid-hostname:6650",
+       })
+
+       assertNil(t, err)
+
+       defer client.Close()
+
+       consumer, err := client.Subscribe(ConsumerOptions{
+               Topic:            "my-topic",
+               SubscriptionName: "my-subscription",
+       })
+
+       // Expect error in creating consumer
+       assertNil(t, consumer)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), ConnectError);
+}
+
+func TestConsumer(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://localhost:6650",
+       })
+
+       assertNil(t, err)
+       defer client.Close()
+
+       producer, err := client.CreateProducer(ProducerOptions{
+               Topic: "my-topic",
+       })
+
+       assertNil(t, err)
+       defer producer.Close()
+
+       consumer, err := client.Subscribe(ConsumerOptions{
+               Topic:                                     "my-topic",
+               SubscriptionName:                          "my-sub",
+               AckTimeout:                                1 * time.Minute,
+               Name:                                      "my-consumer-name",
+               ReceiverQueueSize:                         100,
+               MaxTotalReceiverQueueSizeAcrossPartitions: 10000,
+               Type:                                      Shared,
+       })
+
+       assertNil(t, err)
+       defer consumer.Close()
+
+       assertEqual(t, consumer.Topic(), "persistent://public/default/my-topic")
+       assertEqual(t, consumer.Subscription(), "my-sub")
+
+       ctx := context.Background()
+
+       for i := 0; i < 10; i++ {
+               if err := producer.Send(ctx, ProducerMessage{
+                       Payload: []byte(fmt.Sprintf("hello-%d", i)),
+               }); err != nil {
+                       t.Fatal(err)
+               }
+
+               msg, err := consumer.Receive(ctx)
+               assertNil(t, err)
+               assertNotNil(t, msg)
+
+               assertEqual(t, string(msg.Payload()), fmt.Sprintf("hello-%d", 
i))
+
+               consumer.Ack(msg)
+       }
+
+       consumer.Unsubscribe()
+}
+
+func TestConsumerWithInvalidConf(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://localhost:6650",
+       })
+
+       if err != nil {
+               t.Fatal(err)
+               return
+       }
+
+       defer client.Close()
+
+       consumer, err := client.Subscribe(ConsumerOptions{
+               Topic: "my-topic",
+       })
+
+       // Expect error in creating cosnumer
+       assertNil(t, consumer)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), InvalidConfiguration)
+
+       consumer, err = client.Subscribe(ConsumerOptions{
+               SubscriptionName: "my-subscription",
+       })
+
+       // Expect error in creating cosnumer
+       assertNil(t, consumer)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), InvalidConfiguration)
+}
diff --git a/pulsar-client-go/pulsar/error.go b/pulsar-client-go/pulsar/error.go
new file mode 100644
index 0000000000..a4986c2323
--- /dev/null
+++ b/pulsar-client-go/pulsar/error.go
@@ -0,0 +1,58 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+type Result int
+
+const (
+       UnknownError                          Result = 1  // Unknown error 
happened on broker
+       InvalidConfiguration                  Result = 2  // Invalid 
configuration
+       TimeoutError                          Result = 3  // Operation timed out
+       LookupError                           Result = 4  // Broker lookup 
failed
+       ConnectError                          Result = 5  // Failed to connect 
to broker
+       ReadError                             Result = 6  // Failed to read 
from socket
+       AuthenticationError                   Result = 7  // Authentication 
failed on broker
+       AuthorizationError                    Result = 8  // Client is not 
authorized to create producer/consumer
+       ErrorGettingAuthenticationData        Result = 9  // Client cannot find 
authorization data
+       BrokerMetadataError                   Result = 10 // Broker failed in 
updating metadata
+       BrokerPersistenceError                Result = 11 // Broker failed to 
persist entry
+       ChecksumError                         Result = 12 // Corrupt message 
checksum failure
+       ConsumerBusy                          Result = 13 // Exclusive consumer 
is already connected
+       NotConnectedError                     Result = 14 // Producer/Consumer 
is not currently connected to broker
+       AlreadyClosedError                    Result = 15 // Producer/Consumer 
is already closed and not accepting any operation
+       InvalidMessage                        Result = 16 // Error in 
publishing an already used message
+       ConsumerNotInitialized                Result = 17 // Consumer is not 
initialized
+       ProducerNotInitialized                Result = 18 // Producer is not 
initialized
+       TooManyLookupRequestException         Result = 19 // Too Many 
concurrent LookupRequest
+       InvalidTopicName                      Result = 20 // Invalid topic name
+       InvalidUrl                            Result = 21 // Client Initialized 
with Invalid Broker Url (VIP Url passed to Client Constructor)
+       ServiceUnitNotReady                   Result = 22 // Service Unit 
unloaded between client did lookup and producer/consumer got created
+       OperationNotSupported                 Result = 23
+       ProducerBlockedQuotaExceededError     Result = 24 // Producer is blocked
+       ProducerBlockedQuotaExceededException Result = 25 // Producer is 
getting exception
+       ProducerQueueIsFull                   Result = 26 // Producer queue is 
full
+       MessageTooBig                         Result = 27 // Trying to send a 
messages exceeding the max size
+       TopicNotFound                         Result = 28 // Topic not found
+       SubscriptionNotFound                  Result = 29 // Subscription not 
found
+       ConsumerNotFound                      Result = 30 // Consumer not found
+       UnsupportedVersionError               Result = 31 // Error when an 
older client/version doesn't support a required feature
+       TopicTerminated                       Result = 32 // Topic was already 
terminated
+       CryptoError                           Result = 33 // Error when crypto 
operation fails
+)
diff --git a/pulsar-client-go/pulsar/logger.go 
b/pulsar-client-go/pulsar/logger.go
new file mode 100644
index 0000000000..79734c49c7
--- /dev/null
+++ b/pulsar-client-go/pulsar/logger.go
@@ -0,0 +1,47 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import "fmt"
+
+type LoggerLevel int
+
+const (
+       DEBUG LoggerLevel = iota
+       INFO
+       WARN
+       ERROR
+)
+
+func (l LoggerLevel) String() string {
+       switch l {
+       case DEBUG:
+               return "DEBUG"
+       case INFO:
+               return "INFO"
+       case WARN:
+               return "WARN"
+       case ERROR:
+               return "ERROR"
+
+       default:
+               return fmt.Sprintf("UNKNOWN: %d", l)
+       }
+}
\ No newline at end of file
diff --git a/pulsar-client-go/pulsar/message.go 
b/pulsar-client-go/pulsar/message.go
new file mode 100644
index 0000000000..9b05e37aca
--- /dev/null
+++ b/pulsar-client-go/pulsar/message.go
@@ -0,0 +1,82 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import "time"
+
+type ProducerMessage struct {
+       // Payload for the message
+       Payload []byte
+
+       // Sets the key of the message for routing policy
+       Key string
+
+       // Attach application defined properties on the message
+       Properties map[string]string
+
+       // Set the event time for a given message
+       EventTime time.Time
+
+       // Override the replication clusters for this message.
+       ReplicationClusters []string
+}
+
+type Message interface {
+       // Return the properties attached to the message.
+       // Properties are application defined key/value pairs that will be 
attached to the message
+       Properties() map[string]string
+
+       // Get the payload of the message
+       Payload() []byte
+
+       // Get the unique message ID associated with this message.
+       // The message id can be used to univocally refer to a message without 
having the keep the entire payload in memory.
+       ID() MessageID
+
+       // Get the publish time of this message. The publish time is the 
timestamp that a client publish the message.
+       PublishTime() time.Time
+
+       // Get the event time associated with this message. It is typically set 
by the applications via
+       // `ProducerMessage.EventTime`.
+       // If there isn't any event time associated with this event, it will be 
nil.
+       EventTime() *time.Time
+
+       // Get the key of the message, if any
+       Key() string
+}
+
+// Identifier for a particular message
+type MessageID interface {
+       // Serialize the message id into a sequence of bytes that can be stored 
somewhere else
+       Serialize() []byte
+}
+
+// Reconstruct a MessageID object from its serialized representation
+func DeserializeMessageID(data []byte) MessageID {
+       return deserializeMessageId(data)
+}
+
+var (
+       // MessageID that points to the earliest message avaialable in a topic
+       EarliestMessage MessageID = earliestMessageID()
+
+       // MessageID that points to the latest message
+       LatestMessage MessageID = latestMessageID()
+)
diff --git a/pulsar-client-go/pulsar/pointer.go 
b/pulsar-client-go/pulsar/pointer.go
new file mode 100644
index 0000000000..3e50119214
--- /dev/null
+++ b/pulsar-client-go/pulsar/pointer.go
@@ -0,0 +1,64 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+// #include <stdlib.h>
+import "C"
+import (
+       "unsafe"
+       "sync"
+)
+
+// Inspired by https://github.com/mattn/go-pointer
+// Make sure the marker pointer is freed after restoring
+
+var (
+       mutex sync.Mutex
+       pointers = map[unsafe.Pointer]interface{}{}
+)
+
+func savePointer(object interface{}) unsafe.Pointer {
+       // Get a ref to object using reflection
+       ptr := C.malloc(C.size_t(1))
+
+       mutex.Lock()
+       pointers[ptr] = object
+       mutex.Unlock()
+
+       return ptr
+}
+
+func restorePointer(ptr unsafe.Pointer) interface{} {
+       mutex.Lock()
+       obj := pointers[ptr]
+       delete(pointers, ptr)
+       C.free(ptr)
+       mutex.Unlock()
+
+       return obj
+}
+
+func restorePointerNoDelete(ptr unsafe.Pointer) interface{} {
+       mutex.Lock()
+       obj := pointers[ptr]
+       mutex.Unlock()
+
+       return obj
+}
diff --git a/pulsar-client-go/pulsar/producer.go 
b/pulsar-client-go/pulsar/producer.go
new file mode 100644
index 0000000000..2cfd141633
--- /dev/null
+++ b/pulsar-client-go/pulsar/producer.go
@@ -0,0 +1,166 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import (
+       "time"
+       "context"
+)
+
+type MessageRoutingMode int
+
+const (
+       // Publish messages across all partitions in round-robin.
+       RoundRobinDistribution MessageRoutingMode = 0
+
+       // The producer will chose one single partition and publish all the 
messages into that partition
+       UseSinglePartition MessageRoutingMode = 1
+
+       // Use custom message router implementation that will be called to 
determine the partition for a particular message.
+       CustomPartition MessageRoutingMode = 2
+)
+
+type HashingScheme int
+
+const (
+       JavaStringHash HashingScheme = 0 // Java String.hashCode() equivalent
+       Murmur3_32Hash HashingScheme = 1 // Use Murmur3 hashing function
+       BoostHash      HashingScheme = 2 // C++ based boost::hash
+)
+
+type CompressionType int
+
+const (
+       NoCompression CompressionType = 0
+       LZ4           CompressionType = 1
+       ZLib          CompressionType = 2
+)
+
+type TopicMetadata interface {
+       // Get the number of partitions for the specific topic
+       NumPartitions() int
+}
+
+type ProducerOptions struct {
+       // Specify the topic this producer will be publishing on.
+       // This argument is required when constructing the producer.
+       Topic string
+
+       // Specify a name for the producer
+       // If not assigned, the system will generate a globally unique name 
which can be access with
+       // Producer.ProducerName().
+       // When specifying a name, it is up to the user to ensure that, for a 
given topic, the producer name is unique
+       // across all Pulsar's clusters. Brokers will enforce that only a 
single producer a given name can be publishing on
+       // a topic.
+       Name string
+
+       // Set the send timeout (default: 30 seconds)
+       // If a message is not acknowledged by the server before the 
sendTimeout expires, an error will be reported.
+       // Setting the timeout to -1, will set the timeout to infinity, which 
can be useful when using Pulsar's message
+       // deduplication feature.
+       SendTimeout time.Duration
+
+       // Set the max size of the queue holding the messages pending to 
receive an acknowledgment from the broker.
+       // When the queue is full, by default, all calls to Producer.send() and 
Producer.sendAsync() will fail
+       // unless `BlockIfQueueFull` is set to true. Use 
BlockIfQueueFull(boolean) to change the blocking behavior.
+       MaxPendingMessages int
+
+       // Set the number of max pending messages across all the partitions
+       // This setting will be used to lower the max pending messages for each 
partition
+       // `MaxPendingMessages(int)`, if the total exceeds the configured value.
+       MaxPendingMessagesAcrossPartitions int
+
+       // Set whether the `Producer.Send()` and `Producer.sendAsync()` 
operations should block when the outgoing
+       // message queue is full. Default is `false`. If set to `false`, send 
operations will immediately fail with
+       // `ProducerQueueIsFullError` when there is no space left in pending 
queue.
+       BlockIfQueueFull bool
+
+       // Set the message routing mode for the partitioned producer.
+       // Default routing mode is round-robin routing.
+       //
+       // This logic is applied when the application is not setting a key 
ProducerMessage#setKey(String) on a
+       // particular message.
+       MessageRoutingMode
+
+       // Change the `HashingScheme` used to chose the partition on where to 
publish a particular message.
+       // Standard hashing functions available are:
+       //
+       //  - `JavaStringHash` : Java String.hashCode() equivalent
+       //  - `Murmur3_32Hash` : Use Murmur3 hashing function.
+       //              
https://en.wikipedia.org/wiki/MurmurHash";>https://en.wikipedia.org/wiki/MurmurHash
+       //  - `BoostHash`      : C++ based boost::hash
+       //
+       // Default is `JavaStringHash`.
+       HashingScheme
+
+       // Set the compression type for the producer.
+       // By default, message payloads are not compressed. Supported 
compression types are:
+       //  - LZ4
+       //  - ZLIB
+       CompressionType
+
+       // Set a custom message routing policy by passing an implementation of 
MessageRouter
+       // The router is a function that given a particular message and the 
topic metadata, returns the
+       // partition index where the message should be routed to
+       MessageRouter func(Message, TopicMetadata) int
+
+       // Control whether automatic batching of messages is enabled for the 
producer. Default: false [No batching]
+       //
+       // When batching is enabled, multiple calls to Producer.sendAsync can 
result in a single batch to be sent to the
+       // broker, leading to better throughput, especially when publishing 
small messages. If compression is enabled,
+       // messages will be compressed at the batch level, leading to a much 
better compression ratio for similar headers or
+       // contents.
+       //
+       // When enabled default batch delay is set to 1 ms and default batch 
size is 1000 messages
+       Batching bool
+
+       // Set the time period within which the messages sent will be batched 
(default: 10ms) if batch messages are
+       // enabled. If set to a non zero value, messages will be queued until 
this time interval or until
+       BatchingMaxPublishDelay time.Duration
+
+       // Set the maximum number of messages permitted in a batch. (default: 
1000) If set to a value greater than 1,
+       // messages will be queued until this threshold is reached or batch 
interval has elapsed
+       BatchingMaxMessages uint
+}
+
+// The producer is used to publish messages on a topic
+type Producer interface {
+       // return the topic to which producer is publishing to
+       Topic() string
+
+       // return the producer name which could have been assigned by the 
system or specified by the client
+       Name() string
+
+       // Send a message
+       // This call will be blocking until is successfully acknowledged by the 
Pulsar broker.
+       // Example:
+       // producer.Send(ctx, pulsar.ProducerMessage{ Payload: myPayload })
+       Send(context.Context, ProducerMessage) error
+
+       // Send a message in asynchronous mode
+       // The callback will report back the message being published and
+       // the eventual error in publishing
+       SendAsync(context.Context, ProducerMessage, func(ProducerMessage, 
error))
+
+       // Close the producer and releases resources allocated
+       // No more writes will be accepted from this producer. Waits until all 
pending write request are persisted. In case
+       // of errors, pending writes will not be retried.
+       Close() error
+}
diff --git a/pulsar-client-go/pulsar/producer_test.go 
b/pulsar-client-go/pulsar/producer_test.go
new file mode 100644
index 0000000000..d7748f7bc7
--- /dev/null
+++ b/pulsar-client-go/pulsar/producer_test.go
@@ -0,0 +1,168 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import (
+       "testing"
+       "fmt"
+       "context"
+       "time"
+)
+
+func TestInvalidURL(t *testing.T) {
+       client, err := NewClient(ClientOptions{})
+
+       if client != nil || err == nil {
+               t.Fatal("Should have failed to create client")
+       }
+}
+
+func TestProducerConnectError(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://invalid-hostname:6650",
+       })
+
+       assertNil(t, err)
+
+       defer client.Close()
+
+       producer, err := client.CreateProducer(ProducerOptions{
+               Topic: "my-topic",
+       })
+
+       // Expect error in creating producer
+       assertNil(t, producer)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), ConnectError);
+}
+
+func TestProducer(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL:                      "pulsar://localhost:6650",
+               StatsIntervalInSeconds:   10,
+               IOThreads:                1,
+               OperationTimeoutSeconds:  30,
+               ConcurrentLookupRequests: 1000,
+               MessageListenerThreads:   5,
+               EnableTLS:                false,
+       })
+
+       assertNil(t, err)
+       defer client.Close()
+
+       producer, err := client.CreateProducer(ProducerOptions{
+               Topic:                   "my-topic",
+               Name:                    "my-producer-name",
+               SendTimeout:             10 * time.Second,
+               Batching:                true,
+               BatchingMaxMessages:     100,
+               BatchingMaxPublishDelay: 10 * time.Millisecond,
+               MaxPendingMessages:      100,
+               BlockIfQueueFull:        true,
+               CompressionType:         LZ4,
+       })
+
+       assertNil(t, err)
+       defer producer.Close()
+
+       assertEqual(t, producer.Topic(), "persistent://public/default/my-topic")
+       assertEqual(t, producer.Name(), "my-producer-name")
+
+       ctx := context.Background()
+
+       for i := 0; i < 10; i++ {
+               if err := producer.Send(ctx, ProducerMessage{
+                       Payload: []byte(fmt.Sprintf("hello-%d", i)),
+               }); err != nil {
+                       t.Fatal(err)
+               }
+       }
+}
+
+func TestProducerNoTopic(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://localhost:6650",
+       })
+
+       if err != nil {
+               t.Fatal(err)
+               return
+       }
+
+       defer client.Close()
+
+       producer, err := client.CreateProducer(ProducerOptions{
+       })
+
+       // Expect error in creating producer
+       assertNil(t, producer)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), InvalidConfiguration)
+}
+
+func TestMessageRouter(t *testing.T) {
+       // Create topic with 5 partitions
+       
httpPut("http://localhost:8080/admin/v2/persistent/public/default/my-partitioned-topic/partitions";,
+               5)
+
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://localhost:6650",
+       })
+
+       assertNil(t, err)
+       defer client.Close()
+
+       // Only subscribe on the specific partition
+       consumer, err := client.Subscribe(ConsumerOptions{
+               Topic:            "my-partitioned-topic-partition-2",
+               SubscriptionName: "my-sub",
+       })
+
+       assertNil(t, err)
+       defer consumer.Close()
+
+       producer, err := client.CreateProducer(ProducerOptions{
+               Topic: "my-partitioned-topic",
+               MessageRouter: func(msg Message, tm TopicMetadata) int {
+                       fmt.Println("Routing message ", msg, " -- Partitions: 
", tm.NumPartitions())
+                       return 2
+               },
+       })
+
+       assertNil(t, err)
+       defer producer.Close()
+
+       ctx := context.Background()
+
+       err = producer.Send(ctx, ProducerMessage{
+               Payload: []byte("hello"),
+       })
+       assertNil(t, err)
+
+       fmt.Println("PUBLISHED")
+
+       // Verify message was published on partition 2
+       msg, err := consumer.Receive(ctx)
+       assertNil(t, err)
+       assertNotNil(t, msg)
+       assertEqual(t, string(msg.Payload()), "hello")
+}
diff --git a/pulsar-client-go/pulsar/reader.go 
b/pulsar-client-go/pulsar/reader.go
new file mode 100644
index 0000000000..f61ebd7410
--- /dev/null
+++ b/pulsar-client-go/pulsar/reader.go
@@ -0,0 +1,72 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import "context"
+
+type ReaderMessage struct {
+       Reader
+       Message
+}
+
+type ReaderOptions struct {
+       // Specify the topic this consumer will subscribe on.
+       // This argument is required when constructing the reader.
+       Topic string
+
+       // Set the reader name.
+       Name string
+
+       // The initial reader positioning is done by specifying a message id. 
The options are:
+       //  * `pulsar.EarliestMessage` : Start reading from the earliest 
message available in the topic
+       //  * `pulsar.LatestMessage` : Start reading from the end topic, only 
getting messages published after the
+       //                           reader was created
+       //  * `MessageID` : Start reading from a particular message id, the 
reader will position itself on that
+       //                  specific position. The first message to be read 
will be the message next to the specified
+       //                  messageID
+       StartMessageID MessageID
+
+       // Sets a `MessageChannel` for the consumer
+       // When a message is received, it will be pushed to the channel for 
consumption
+       MessageChannel chan ReaderMessage
+
+       // Sets the size of the consumer receive queue.
+       // The consumer receive queue controls how many messages can be 
accumulated by the Reader before the
+       // application calls Reader.readNext(). Using a higher value could 
potentially increase the consumer
+       // throughput at the expense of bigger memory utilization.
+       //
+       // Default value is {@code 1000} messages and should be good for most 
use cases.
+       ReceiverQueueSize int
+
+       // Set the subscription role prefix. The default prefix is "reader".
+       SubscriptionRolePrefix string
+}
+
+// A Reader can be used to scan through all the messages currently available 
in a topic.
+type Reader interface {
+       // The topic from which this reader is reading from
+       Topic() string
+
+       // Read the next message in the topic, blocking until a message is 
available
+       Next(context.Context) (Message, error)
+
+       // Close the reader and stop the broker to push more messages
+       Close() error
+}
diff --git a/pulsar-client-go/pulsar/reader_test.go 
b/pulsar-client-go/pulsar/reader_test.go
new file mode 100644
index 0000000000..11d1b3620c
--- /dev/null
+++ b/pulsar-client-go/pulsar/reader_test.go
@@ -0,0 +1,122 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import (
+       "testing"
+       "fmt"
+       "context"
+)
+
+func TestReaderConnectError(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://invalid-hostname:6650",
+       })
+
+       assertNil(t, err)
+
+       defer client.Close()
+
+       reader, err := client.CreateReader(ReaderOptions{
+               Topic:          "my-topic",
+               StartMessageID: EarliestMessage,
+       })
+
+       // Expect error in creating reader
+       assertNil(t, reader)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), ConnectError);
+}
+
+func TestReader(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://localhost:6650",
+       })
+
+       assertNil(t, err)
+       defer client.Close()
+
+       producer, err := client.CreateProducer(ProducerOptions{
+               Topic: "my-reader-topic",
+       })
+
+       assertNil(t, err)
+       defer producer.Close()
+
+       reader, err := client.CreateReader(ReaderOptions{
+               Topic:          "my-reader-topic",
+               StartMessageID: LatestMessage,
+       })
+
+       assertNil(t, err)
+       defer reader.Close()
+
+       assertEqual(t, reader.Topic(), 
"persistent://public/default/my-reader-topic")
+
+       ctx := context.Background()
+
+       for i := 0; i < 10; i++ {
+               if err := producer.Send(ctx, ProducerMessage{
+                       Payload: []byte(fmt.Sprintf("hello-%d", i)),
+               }); err != nil {
+                       t.Fatal(err)
+               }
+
+               msg, err := reader.Next(ctx)
+               assertNil(t, err)
+               assertNotNil(t, msg)
+
+               assertEqual(t, string(msg.Payload()), fmt.Sprintf("hello-%d", 
i))
+       }
+}
+
+func TestReaderWithInvalidConf(t *testing.T) {
+       client, err := NewClient(ClientOptions{
+               URL: "pulsar://localhost:6650",
+       })
+
+       if err != nil {
+               t.Fatal(err)
+               return
+       }
+
+       defer client.Close()
+
+       reader, err := client.CreateReader(ReaderOptions{
+               Topic: "my-topic",
+       })
+
+       // Expect error in creating cosnumer
+       assertNil(t, reader)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), InvalidConfiguration)
+
+       reader, err = client.CreateReader(ReaderOptions{
+               StartMessageID: LatestMessage,
+       })
+
+       // Expect error in creating cosnumer
+       assertNil(t, reader)
+       assertNotNil(t, err)
+
+       assertEqual(t, err.(*Error).Result(), InvalidConfiguration)
+}
diff --git a/pulsar-client-go/pulsar/util_test.go 
b/pulsar-client-go/pulsar/util_test.go
new file mode 100644
index 0000000000..5fdbac2bd3
--- /dev/null
+++ b/pulsar-client-go/pulsar/util_test.go
@@ -0,0 +1,69 @@
+//
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+//
+
+package pulsar
+
+import (
+       "testing"
+       "runtime"
+       "net/http"
+       "log"
+       "encoding/json"
+       "bytes"
+)
+
+func assertNil(t *testing.T, a interface{}) {
+       if a != nil {
+               _, file, line, _ := runtime.Caller(1)
+               t.Fatalf("%s:%d  | Expected nil", file, line)
+       }
+}
+
+func assertNotNil(t *testing.T, a interface{}) {
+       if a == nil {
+               _, file, line, _ := runtime.Caller(1)
+               t.Fatalf("%s:%d  | Expected not nil", file, line)
+       }
+}
+
+func assertEqual(t *testing.T, realValue interface{}, expected interface{}) {
+       if realValue != expected {
+               _, file, line, _ := runtime.Caller(1)
+               t.Fatalf("%s:%d  | Expected '%v' -- Got '%v'", file, line, 
expected, realValue)
+       }
+}
+
+func httpPut(url string, body interface{}) {
+       client := http.Client{}
+
+       data, _ := json.Marshal(body)
+       req, err := http.NewRequest(http.MethodPut, url, bytes.NewReader(data))
+       if err != nil {
+               log.Fatal(err)
+       }
+
+       req.Header = map[string][]string{
+               "Content-Type": {"application/json"},
+       }
+
+       _, err = client.Do(req)
+       if err != nil {
+               log.Fatal(err)
+       }
+}


 

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
[email protected]


With regards,
Apache Git Services

Reply via email to