diff --git a/LICENSE b/LICENSE
new file mode 100644
--- /dev/null
+++ b/LICENSE
@@ -0,0 +1,20 @@
+Copyright (c) 2016 Alexey Raga
+
+Permission is hereby granted, free of charge, to any person obtaining
+a copy of this software and associated documentation files (the
+"Software"), to deal in the Software without restriction, including
+without limitation the rights to use, copy, modify, merge, publish,
+distribute, sublicense, and/or sell copies of the Software, and to
+permit persons to whom the Software is furnished to do so, subject to
+the following conditions:
+
+The above copyright notice and this permission notice shall be included
+in all copies or substantial portions of the Software.
+
+THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND,
+EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF
+MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT.
+IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY
+CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT,
+TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE
+SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
diff --git a/Setup.hs b/Setup.hs
new file mode 100644
--- /dev/null
+++ b/Setup.hs
@@ -0,0 +1,2 @@
+import Distribution.Simple
+main = defaultMain
diff --git a/example/Main.hs b/example/Main.hs
new file mode 100644
--- /dev/null
+++ b/example/Main.hs
@@ -0,0 +1,49 @@
+{-# LANGUAGE OverloadedStrings #-}
+module Main where
+
+import Control.Monad (forM_)
+import Data.Monoid ((<>))
+import Kafka.Conduit.Source
+import Kafka.Conduit.Sink
+import Data.Conduit
+import qualified Data.Conduit.List as L
+
+kafkaBroker :: BrokerAddress
+kafkaBroker = BrokerAddress "localhost:9092"
+
+-- Topic to write to and read from
+testTopic :: TopicName
+testTopic = TopicName "kafka-client-conduit-example-topic"
+
+-- Global consumer properties
+consumerProps :: ConsumerProperties
+consumerProps = consumerBrokersList [kafkaBroker]
+             <> groupId (ConsumerGroupId "consumer_conduit_example_group")
+             <> noAutoCommit
+
+-- Subscription to topics
+consumerSub :: Subscription
+consumerSub = topics [testTopic]
+           <> offsetReset Earliest
+
+-- Global producer properties
+producerProps :: ProducerProperties
+producerProps = producerBrokersList [kafkaBroker]
+
+main :: IO ()
+main = do
+  putStrLn "Running sink..."
+  producerRes <- runConduitRes outputStream
+  forM_ producerRes print
+  putStrLn "Running source..."
+  msgs <- runConduitRes $ inputStream .| L.take 10
+  forM_ msgs print
+  putStrLn "Ok."
+
+inputStream =
+  kafkaSource consumerProps consumerSub (Timeout 1000)
+
+outputStream =
+  L.sourceList ["c-one", "c-two", "c-three"]
+  .| L.map ((ProducerRecord testTopic UnassignedPartition Nothing) . Just)
+  .| kafkaSink producerProps
diff --git a/hw-kafka-conduit.cabal b/hw-kafka-conduit.cabal
new file mode 100644
--- /dev/null
+++ b/hw-kafka-conduit.cabal
@@ -0,0 +1,71 @@
+name:                hw-kafka-conduit
+version:             1.0.0
+synopsis:            Conduit bindings for kafka-client
+homepage:            https://github.com/haskell-works/hw-kafka-client-conduit
+bug-reports:         https://github.com/haskell-works/hw-kafka-client-conduit/issues
+license:             MIT
+license-file:        LICENSE
+author:              Alexey Raga <alexey.raga@gmail.com>
+maintainer:          Alexey Raga <alexey.raga@gmail.com>
+category:            Database
+copyright:           Alexey Raga
+build-type:          Simple
+-- extra-source-files:
+cabal-version:       >=1.10
+description:         Conduit bindings for hw-kafka-client
+
+source-repository head
+  type:     git
+  location: https://github.com/haskell-works/hw-kafka-conduit
+
+executable kafka-conduit-example
+  main-is:             Main.hs
+  hs-source-dirs:      example
+  default-language:    Haskell2010
+  build-depends:       base >=4.6 && < 5
+                     , bifunctors
+                     , bytestring
+                     , conduit
+                     , containers
+                     , resourcet
+                     , hw-kafka-conduit
+
+library
+  hs-source-dirs:      src
+  ghc-options:         -Wall
+  exposed-modules:     Kafka.Conduit
+                     , Kafka.Conduit.Source
+                     , Kafka.Conduit.Sink
+                     , Kafka.Conduit.Utils
+  default-language:    Haskell2010
+  build-depends:       base >= 4.7 && < 5
+                     , bifunctors
+                     , bytestring
+                     , containers
+                     , conduit
+                     , conduit-extra
+                     , exceptions
+                     , hw-kafka-client
+                     , mtl
+                     , resourcet
+                     , transformers
+
+
+test-suite kafka-client-conduit-test
+  type:                exitcode-stdio-1.0
+  hs-source-dirs:      test
+  main-is:             Spec.hs
+  ghc-options:         -Wall -threaded -rtsopts -with-rtsopts=-N
+  default-language:    Haskell2010
+  build-depends:       base
+                     , hw-kafka-conduit
+                     , bifunctors
+                     , bytestring
+                     , containers
+                     , conduit
+                     , conduit-extra
+                     , hspec
+                     , hw-kafka-client
+                     , mtl
+                     , resourcet
+                     , transformers
diff --git a/src/Kafka/Conduit.hs b/src/Kafka/Conduit.hs
new file mode 100644
--- /dev/null
+++ b/src/Kafka/Conduit.hs
@@ -0,0 +1,6 @@
+module Kafka.Conduit
+( module X
+) where
+
+import Kafka.Conduit.Source as X
+import Kafka.Conduit.Sink as X
diff --git a/src/Kafka/Conduit/Sink.hs b/src/Kafka/Conduit/Sink.hs
new file mode 100644
--- /dev/null
+++ b/src/Kafka/Conduit/Sink.hs
@@ -0,0 +1,63 @@
+module Kafka.Conduit.Sink
+( module X
+, kafkaSink, kafkaSinkAutoClose, kafkaSinkNoClose
+) where
+--
+
+import Control.Monad.IO.Class
+import Control.Monad (void)
+import Control.Monad.Trans.Resource
+import Data.Conduit
+import Kafka.Producer as X
+import Kafka.Conduit.Utils as X
+
+kafkaSinkAutoClose :: MonadResource m
+                   => KafkaProducer
+                   -> Sink ProducerRecord m (Maybe KafkaError)
+kafkaSinkAutoClose prod =
+  bracketP (return prod) (void . closeProducer) runHandler
+  where
+    runHandler p' = do
+      mbMsg <- await
+      case mbMsg of
+        Nothing -> return Nothing
+        Just msg -> do
+          res <- produceMessage p' msg
+          case res of
+            Nothing -> runHandler p'
+            Just err -> return (Just err)
+
+kafkaSinkNoClose :: MonadIO m
+                 => KafkaProducer
+                 -> Sink ProducerRecord m (Maybe KafkaError)
+kafkaSinkNoClose prod = go
+  where
+    go = do
+      mbMsg <- await
+      case mbMsg of
+        Nothing -> return Nothing
+        Just msg -> do
+          res <- produceMessage prod msg
+          case res of
+            Nothing -> go
+            Just err -> return (Just err)
+
+kafkaSink :: MonadResource m
+          => ProducerProperties
+          -> Sink ProducerRecord m (Maybe KafkaError)
+kafkaSink props =
+  bracketP mkProducer clProducer runHandler
+  where
+    mkProducer = newProducer props
+
+    clProducer (Left _) = return ()
+    clProducer (Right prod) = void $ closeProducer prod
+
+    runHandler (Left err) = return (Just err)
+    runHandler (Right prod) = do
+      mbMsg <- await
+      case mbMsg of
+        Nothing -> return Nothing
+        Just msg -> do
+          res <- produceMessage prod msg
+          runHandler $ maybe (Right prod) Left res
diff --git a/src/Kafka/Conduit/Source.hs b/src/Kafka/Conduit/Source.hs
new file mode 100644
--- /dev/null
+++ b/src/Kafka/Conduit/Source.hs
@@ -0,0 +1,168 @@
+{-# LANGUAGE TupleSections#-}
+module Kafka.Conduit.Source
+( module X
+, kafkaSource, kafkaSourceNoClose, kafkaSourceAutoClose
+, mapFirst, mapF, bimapF
+, sequenceValueFirst, sequenceValue, bisequenceValue
+, traverseValueFirst, traverseValue, bitraverseValue
+, traverseValueFirstM, traverseValueM, bitraverseValueM
+, isFatal, isPollTimeout, isPartitionEOF
+, skipNonFatal, nonFatalOr
+) where
+
+import Data.Bifunctor
+import Data.Bitraversable
+import Control.Monad.IO.Class
+import Control.Monad (void)
+import Control.Monad.Trans.Resource
+import qualified Data.ByteString as BS
+import Data.Conduit
+import qualified Data.Conduit.List as L
+import Kafka.Consumer as X
+import Kafka.Conduit.Utils as X
+
+kafkaSourceNoClose :: MonadIO m
+                   => KafkaConsumer
+                   -> Timeout
+                   -> Source m (Either KafkaError (ConsumerRecord (Maybe BS.ByteString) (Maybe BS.ByteString)))
+kafkaSourceNoClose c t = go
+  where
+    go = do
+      msg <- pollMessage c t
+      -- stop at some certain cases because it is not goind to be better with time
+      case msg of
+        Left err | isFatal err -> void $ yield (Left err)
+        _ -> yield msg >> go
+
+kafkaSourceAutoClose :: MonadResource m
+                     => KafkaConsumer
+                     -> Timeout
+                     -> Source m (Either KafkaError (ConsumerRecord (Maybe BS.ByteString) (Maybe BS.ByteString)))
+kafkaSourceAutoClose c ts =
+  bracketP mkConsumer clConsumer runHandler
+  where
+    mkConsumer = return c
+    clConsumer c' = void $ closeConsumer c'
+    runHandler c' = do
+      msg <- pollMessage c' ts
+      -- stop at some certain cases because it is not goind to be better with time
+      case msg of
+        Left err | isFatal err -> void $ yield (Left err)
+        _ -> yield msg >> runHandler c'
+
+-- | Conduit source for kafka consumer.
+kafkaSource :: MonadResource m
+            => ConsumerProperties
+            -> Subscription
+            -> Timeout
+            -> Source m (Either KafkaError (ConsumerRecord (Maybe BS.ByteString) (Maybe BS.ByteString)))
+kafkaSource props sub ts =
+  bracketP mkConsumer clConsumer runHandler
+  where
+      mkConsumer = newConsumer props sub
+
+      clConsumer (Left _) = return ()
+      clConsumer (Right c) = void $ closeConsumer c
+
+      runHandler (Left err) = void $ yield (Left err)
+      runHandler (Right c) = do
+        msg <- pollMessage c ts
+        -- stop at some certain cases because it is not goind to be better with time
+        case msg of
+          Left err | isFatal err -> void $ yield (Left err)
+          _ -> yield msg >> runHandler (Right c)
+
+---------------------- Useful ConsumerRecord combinators -----------------------
+mapFirst :: (Bifunctor t, Monad m) => (k -> k') -> Conduit (t k v) m (t k' v)
+mapFirst f = L.map (first f)
+{-# INLINE mapFirst #-}
+
+mapF :: (Functor t, Monad m) => (v -> v') -> Conduit (t v) m (t v')
+mapF f = L.map (fmap f)
+{-# INLINE mapF #-}
+
+bimapF :: (Bifunctor t, Monad m) => (k -> k') -> (v -> v') -> Conduit (t k v) m (t k' v')
+bimapF f g = L.map (bimap f g)
+{-# INLINE bimapF #-}
+
+sequenceValueFirst :: (Bitraversable t, Applicative f, Monad m) => Conduit (t (f k) v) m (f (t k v))
+sequenceValueFirst = L.map sequenceFirst
+{-# INLINE sequenceValueFirst #-}
+
+sequenceValue :: (Traversable t, Applicative f, Monad m) => Conduit (t (f v)) m (f (t v))
+sequenceValue = L.map sequenceA
+{-# INLINE sequenceValue #-}
+
+bisequenceValue :: (Bitraversable t, Applicative f, Monad m) => Conduit (t (f k) (f v)) m (f (t k v))
+bisequenceValue = L.map bisequenceA
+{-# INLINE bisequenceValue #-}
+
+traverseValueFirst :: (Bitraversable t, Applicative f, Monad m) => (k -> f k') -> Conduit (t k v) m (f (t k' v))
+traverseValueFirst f = L.map (traverseFirst f)
+{-# INLINE traverseValueFirst #-}
+
+traverseValue :: (Traversable t, Applicative f, Monad m) => (v -> f v') -> Conduit (t v) m (f (t v'))
+traverseValue f = L.map (traverse f)
+{-# INLINE traverseValue #-}
+
+bitraverseValue :: (Bitraversable t, Applicative f, Monad m) => (k -> f k') -> (v -> f v') -> Conduit (t k v) m (f (t k' v'))
+bitraverseValue f g = L.map (bitraverse f g)
+{-# INLINE bitraverseValue #-}
+
+traverseValueFirstM :: (Bitraversable t, Applicative f, Monad m) => (k -> m (f k')) -> Conduit (t k v) m (f (t k' v))
+traverseValueFirstM f = L.mapM (traverseFirstM f)
+{-# INLINE traverseValueFirstM #-}
+
+traverseValueM :: (Traversable t, Applicative f, Monad m) => (v -> m (f v')) -> Conduit (t v) m (f (t v'))
+traverseValueM f = L.mapM (traverseM f)
+{-# INLINE traverseValueM #-}
+
+bitraverseValueM :: (Bitraversable t, Applicative f, Monad m) => (k -> m (f k')) -> (v -> m (f v')) -> Conduit (t k v) m (f (t k' v'))
+bitraverseValueM f g = L.mapM (bitraverseM f g)
+{-# INLINE bitraverseValueM #-}
+
+--------------------------------------------------------------------------------
+
+skipNonFatal :: Monad m => Conduit (Either KafkaError b) m (Either KafkaError b)
+skipNonFatal = L.filter (either isFatal (const True))
+{-# INLINE skipNonFatal #-}
+
+isPollTimeout :: KafkaError -> Bool
+isPollTimeout e = KafkaResponseError RdKafkaRespErrTimedOut == e
+{-# INLINE isPollTimeout #-}
+
+isPartitionEOF :: KafkaError -> Bool
+isPartitionEOF e = KafkaResponseError RdKafkaRespErrPartitionEof == e
+{-# INLINE isPartitionEOF #-}
+
+nonFatalOr :: Monad m => [KafkaError -> Bool] -> Conduit (Either KafkaError b) m (Either KafkaError b)
+nonFatalOr fs =
+  let fun e = or $ (\f -> f e) <$> (isFatal : fs)
+   in L.filter (either fun (const True))
+{-# INLINE nonFatalOr #-}
+
+-- | Checks if the error is fatal in a way that it doesn't make sense to retry.
+isFatal :: KafkaError -> Bool
+isFatal e = case e of
+  KafkaUnknownConfigurationKey _   -> True
+  KafkaInvalidConfigurationValue _ -> True
+  KakfaBadConfiguration            -> True
+  KafkaBadSpecification _          -> True
+
+  -- More of them? Less of them?
+  KafkaResponseError RdKafkaRespErrDestroy                    -> True
+  KafkaResponseError RdKafkaRespErrFail                       -> True
+  KafkaResponseError RdKafkaRespErrInvalidArg                 -> True
+  KafkaResponseError RdKafkaRespErrSsl                        -> True
+  KafkaResponseError RdKafkaRespErrUnknownProtocol            -> True
+  KafkaResponseError RdKafkaRespErrNotImplemented             -> True
+  KafkaResponseError RdKafkaRespErrAuthentication             -> True
+  KafkaResponseError RdKafkaRespErrInconsistentGroupProtocol  -> True
+  KafkaResponseError RdKafkaRespErrTopicAuthorizationFailed   -> True
+  KafkaResponseError RdKafkaRespErrGroupAuthorizationFailed   -> True
+  KafkaResponseError RdKafkaRespErrClusterAuthorizationFailed -> True
+  KafkaResponseError RdKafkaRespErrUnsupportedSaslMechanism   -> True
+  KafkaResponseError RdKafkaRespErrIllegalSaslState           -> True
+  KafkaResponseError RdKafkaRespErrUnsupportedVersion         -> True
+
+  _ -> False
diff --git a/src/Kafka/Conduit/Utils.hs b/src/Kafka/Conduit/Utils.hs
new file mode 100644
--- /dev/null
+++ b/src/Kafka/Conduit/Utils.hs
@@ -0,0 +1,22 @@
+module Kafka.Conduit.Utils where
+
+import Control.Exception
+import Control.Monad
+import Control.Monad.Catch
+import Data.Conduit
+
+throwLeft :: (MonadThrow m, Exception e) => Conduit (Either e i) m i
+throwLeft = awaitForever (either throwM yield)
+
+-- | Create a conduit that folds with the function f over its input i with its
+-- internal state s and emits outputs [o], then finally emits outputs [o] from
+-- the function g applied to the final state s.
+foldYield :: Monad m => (i -> s -> (s, [o])) -> (s -> [o]) -> s -> Conduit i m o
+foldYield f g s = do
+  mi <- await
+  case mi of
+    Just i -> do
+      let (s', os) = f i s
+      forM_ os yield
+      foldYield f g s'
+    Nothing -> forM_ (g s) yield
diff --git a/test/Spec.hs b/test/Spec.hs
new file mode 100644
--- /dev/null
+++ b/test/Spec.hs
@@ -0,0 +1,2 @@
+main :: IO ()
+main = putStrLn "Test suite not yet implemented"
