diff --git a/cqrs-sqlite3.cabal b/cqrs-sqlite3.cabal
--- a/cqrs-sqlite3.cabal
+++ b/cqrs-sqlite3.cabal
@@ -1,5 +1,5 @@
 Name:                cqrs-sqlite3
-Version:             0.7.1
+Version:             0.8.0
 Synopsis:            SQLite3 backend for the cqrs package.
 Description:         SQLite3 backend for the cqrs package.
 License:             MIT
@@ -14,8 +14,8 @@
   build-depends:   base == 4.*
                  , bytestring >= 0.9.0.1
                  , cereal >= 0.3.3 && < 0.4
-                 , cqrs >= 0.7.1 && < 0.8
-                 , enumerator >= 0.4.15 && < 0.5
+                 , conduit >= 0.1 && < 0.2
+                 , cqrs >= 0.8.0 && < 0.9
                  , direct-sqlite >= 1.1 && < 1.2
   extensions:      ScopedTypeVariables
   ghc-options:     -Wall
diff --git a/src/Data/CQRS/EventStore/Backend/Sqlite3.hs b/src/Data/CQRS/EventStore/Backend/Sqlite3.hs
--- a/src/Data/CQRS/EventStore/Backend/Sqlite3.hs
+++ b/src/Data/CQRS/EventStore/Backend/Sqlite3.hs
@@ -5,14 +5,13 @@
 
 import           Control.Monad (when, forM_, liftM)
 import           Data.ByteString (ByteString)
-import           Data.CQRS.EventStore.Backend (EventStoreBackend(..), RawEvent)
-import           Data.CQRS.EventStore.Backend.Sqlite3Utils (withTransaction, execSql, enumQueryResult)
+import           Data.Conduit (Source, ($=), ($$), runResourceT)
+import qualified Data.Conduit.List as CL
+import           Data.CQRS.EventStore.Backend (EventStoreBackend(..), RawEvent, RawSnapshot(..))
+import           Data.CQRS.EventStore.Backend.Sqlite3Utils (withTransaction, execSql, sourceQuery)
 import           Data.CQRS.GUID (GUID)
 import           Data.CQRS.PersistedEvent (PersistedEvent(..))
-import           Data.CQRS.Serialize (decode')
-import           Data.Enumerator (Enumerator, (>>==), ($=), run_)
-import qualified Data.Enumerator.List as EL
-import           Data.Serialize (encode)
+import           Data.CQRS.Serialize (decode', encode)
 import qualified Database.SQLite3 as SQL
 import           Database.SQLite3 (Database, SQLData(..))
 import           Prelude hiding (catch)
@@ -79,8 +78,7 @@
   let unpackColumns [ SQLInteger v ] = v
       unpackColumns columns = error $ badQueryResultMsg [show guid] columns
   -- Get the current version number of the aggregate.
-  curVer <- run_ $ EL.fold (\x -> max x . unpackColumns) 0 >>==
-              (enumQueryResult database getCurrentVersionSql [toSQLData guid])
+  curVer <- runResourceT $ (sourceQuery database getCurrentVersionSql [toSQLData guid]) $$ CL.fold (\x -> max x . unpackColumns) 0
 
   -- Sanity check current version number.
   when (fromIntegral curVer /= originatingVersion) $
@@ -101,19 +99,17 @@
       , toSQLData $ peGlobalVer e
       ]
 
-retrieveEvents :: Database -> GUID -> Int -> IO [RawEvent]
+retrieveEvents :: Database -> GUID -> Int -> Source IO RawEvent
 retrieveEvents database guid v0 = do
   -- Unpack the columns into tuples.
   let unpackColumns [SQLInteger version, SQLInteger gversion, SQLBlob eventData] = PersistedEvent guid eventData (fromIntegral version) (fromIntegral gversion)
       unpackColumns columns = error $ badQueryResultMsg [show guid, show v0] columns
   -- Find events with version numbers.
-  run_ $ EL.consume >>==
-    (enumQueryResult database selectEventsSql [toSQLData guid, toSQLData v0] $=
-     (EL.map unpackColumns))
+  fmap unpackColumns $ sourceQuery database selectEventsSql [toSQLData guid, toSQLData v0]
 
-enumerateAllEvents :: Database -> Int -> Enumerator RawEvent IO a
+enumerateAllEvents :: Database -> Int -> Source IO RawEvent
 enumerateAllEvents database minVersion = do
-  enumQueryResult database enumerateAllEventsSql [toSQLData minVersion] $= EL.map
+  sourceQuery database enumerateAllEventsSql [toSQLData minVersion] $= CL.map
     (\columns -> do
         case columns of
           [ SQLInteger gv, SQLBlob g, SQLInteger v, SQLBlob ed ] ->
@@ -121,23 +117,21 @@
           _ ->
             error $ badQueryResultMsg [show minVersion] columns)
 
-writeSnapshot :: Database -> GUID -> (Int, ByteString) -> IO ()
-writeSnapshot database guid (v,a) = do
+writeSnapshot :: Database -> GUID -> RawSnapshot -> IO ()
+writeSnapshot database guid (RawSnapshot v d) = do
   execSql database writeSnapshotSql
     [ toSQLData guid
-    , toSQLData a
+    , toSQLData d
     , toSQLData v
     ]
 
-getLatestSnapshot :: Database -> GUID -> IO (Maybe (Int, ByteString))
+getLatestSnapshot :: Database -> GUID -> IO (Maybe RawSnapshot)
 getLatestSnapshot database guid = do
   -- Unpack columns from result.
-  let unpackColumns :: [SQLData] -> Maybe (Int,ByteString)
-      unpackColumns [SQLBlob a, SQLInteger v] = Just (fromIntegral v, a)
-      unpackColumns columns                    = error $ badQueryResultMsg [show guid] columns
+  let unpackColumns [SQLBlob d, SQLInteger v] = Just $ RawSnapshot (fromIntegral v) d
+      unpackColumns columns                   = error $ badQueryResultMsg [show guid] columns
   -- Run the query.
-  run_ $ EL.fold const Nothing >>==
-    (enumQueryResult database selectSnapshotSql [toSQLData guid] $= (EL.map unpackColumns))
+  runResourceT $ (sourceQuery database selectSnapshotSql [toSQLData guid] $= (CL.map unpackColumns)) $$ CL.fold const Nothing
 
 getLatestVersion :: Database -> IO Int
 getLatestVersion database = do
@@ -146,7 +140,7 @@
       unpackColumns [SQLInteger v] = fromIntegral v
       unpackColumns columns        = error $ badQueryResultMsg [] columns
   -- Run the query.
-  liftM head $ run_ $ EL.consume >>== (enumQueryResult database getLatestVersionSql [] $= (EL.map unpackColumns))
+  liftM head $ runResourceT $ (sourceQuery database getLatestVersionSql [] $= (CL.map unpackColumns)) $$ CL.consume
 
 -- | Open an SQLite3-based event store using the named SQLite database file.
 -- The database file is created if it does not exist.
diff --git a/src/Data/CQRS/EventStore/Backend/Sqlite3Utils.hs b/src/Data/CQRS/EventStore/Backend/Sqlite3Utils.hs
--- a/src/Data/CQRS/EventStore/Backend/Sqlite3Utils.hs
+++ b/src/Data/CQRS/EventStore/Backend/Sqlite3Utils.hs
@@ -1,13 +1,15 @@
 {-| Implementation of an SQLite3-based event store. -}
 module Data.CQRS.EventStore.Backend.Sqlite3Utils
-       ( enumQueryResult
+       ( sourceQuery
        , execSql
        , withTransaction
        ) where
 
-import           Control.Exception (catch, bracket, finally, onException, SomeException)
-import           Data.Enumerator (Enumerator, Iteratee(..), tryIO, continue, (>>==), Stream(..))
-import           Data.Enumerator.Internal (checkContinue0)
+import           Control.Exception (catch, bracket, onException, SomeException)
+import           Control.Monad (liftM, when)
+import           Data.Conduit (Source)
+import qualified Data.Conduit as C
+import           Data.IORef (newIORef, readIORef, writeIORef)
 import qualified Database.SQLite3 as SQL
 import           Database.SQLite3 (Database, Statement, SQLData(..), StepResult(..))
 import           Prelude hiding (catch)
@@ -34,26 +36,29 @@
     _ <- SQL.step stmt
     return ()
 
-enumQueryResults :: Statement -> Enumerator [SQLData] IO b
-enumQueryResults stmt = checkContinue0 $ \loop k ->
-  do
-    nextResult <- tryIO $ SQL.step stmt
-    case nextResult of
-      Done -> continue k
-      Row -> do
-        cols <- tryIO $ SQL.columns stmt
-        k (Chunks [cols]) >>== loop
+data State = Unbound
+           | Bound
+             deriving (Eq)
 
--- | Perform a query and enumerate the results. Each result is a list
--- of returned columns.
-enumQueryResult :: Database -> String -> [SQLData] -> Enumerator [SQLData] IO b
-enumQueryResult database sql parameters step = do
-  stmt <- tryIO $ SQL.prepare database sql
-  Iteratee $ finally
-    (do
+sourceQuery :: Database -> String -> [SQLData] -> Source IO [SQLData]
+sourceQuery database sql parameters =
+  C.sourceIO
+  (do
+      stateRef <- newIORef Unbound
+      stmt <- SQL.prepare database sql
+      return (stateRef, stmt))
+  (\(_,stmt) -> SQL.finalize stmt)
+  (\(stateRef,stmt) -> do
+      -- Bind parameters if necessary.
+      state <- readIORef stateRef
+      when (state == Unbound) $ do
         SQL.bind stmt parameters
-        runIteratee $ enumQueryResults stmt step)
-    (SQL.finalize stmt)
+        writeIORef stateRef $ Bound
+      -- Fetch results.
+      nextResult <- SQL.step stmt
+      case nextResult of
+        Done -> return C.Closed
+        Row -> liftM C.Open $ SQL.columns stmt)
 
 -- | Execute an IO action with an active transaction.
 withTransaction :: Database -> IO a -> IO a
